#!/usr/bin/env python3 """Generic Server-Sent-Events 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). STATUS: DOCUMENTED (framing rules from the WHATWG SSE spec, event vocabularies from the OpenAPI specs / discovery document / build-with-claude/streaming.md) · offline-tested on recorded fixtures (tests/shared/fixtures/sse/*) by tests/shared/test_sse_parser.py (2026-09-18 OpenAI/Anthropic; 2026-09-19 xAI/Gemini fixtures cut from live captures in tmp-live/xai/*_stream*.json and tmp-live/gemini-core/*.raw.txt, thoughtSignature blobs truncated). Stdlib only. See docs/architecture/streaming-patterns.md. Layers 1. SSEParser — incremental, byte-oriented, bounded buffer (backpressure friendly): feed(bytes) -> list[SSEEvent] Handles `event:`/`data:`/`id:`/`retry:` fields, comments (`: keep-alive`), multi-line data, CRLF/LF, BOM, chunk boundaries splitting a line, and `data: [DONE]` (Chat Completions sentinel). JSONArrayStreamParser — Gemini's NON-SSE streaming mode (no `?alt=sse`): a pretty-printed JSON array streamed incrementally (`[{` … `}\n,\n{` … `}\n]`); emits each top-level object as soon as it closes. 2. PartialJSON — accumulates argument fragments per tool call and offers a tolerant preview parse. 3. OpenAINormalizer / AnthropicNormalizer / XAINormalizer / GeminiNormalizer — unified vocabulary: text_delta · reasoning_delta · tool_call_start · tool_call_delta · tool_call_done · usage · done · error · raw 4. iter_unified(chunks, provider, framing="sse"|"json_array") — convenience generator over an iterable of byte chunks. Framing differences (the JSON payload is authoritative; `event:` names are absent on two of the four providers): OpenAI event: response.output_text.delta data: {"type":"response.output_text.delta","sequence_number":7,"item_id":…,"delta":"Hi"} Anthropic event: content_block_delta data: {"type":"content_block_delta","index":0,"delta":{"type":"text_delta","text":"Hi"}} xAI chat data: {"object":"chat.completion.chunk","choices":[{"delta":{"reasoning_content":"…"}}]} … data: [DONE] (no event: lines) xAI resp. event: response.output_text.delta (identical to OpenAI incl. sequence_number; `response.reasoning_summary_text.delta` observed) Gemini data: {"candidates":[{"content":{"parts":[{"text":"Hi"}]}}],"usageMetadata":{…}}\r\n\r\n (no event:, no [DONE]; the LAST chunk has finishReason and, on Gemini 3, an EMPTY text part carrying `thoughtSignature`) Gemini Interactions event: step.delta data: {"event_type":"step.delta","index":1,"delta":{"type":"text","text":"OK."}} … event: done / data: [DONE] """ from __future__ import annotations import json from dataclasses import dataclass, field from typing import Any, Iterable, Iterator, Optional # -------------------------------------------------------------------------------------- # 1. Generic SSE parser # -------------------------------------------------------------------------------------- @dataclass class SSEEvent: event: Optional[str] data: str id: Optional[str] = None retry: Optional[int] = None def json(self) -> Any: try: return json.loads(self.data) except json.JSONDecodeError: return None class SSEBufferOverflow(Exception): pass class SSEParser: """Incremental SSE parser. Call feed() with each network chunk (any size); events are returned as soon as their terminating blank line arrives. `max_buffer` bounds memory for a single (possibly hostile) event.""" def __init__(self, max_buffer: int = 8 * 1024 * 1024): self._buf = bytearray() self._event: Optional[str] = None self._data: list[str] = [] self._id: Optional[str] = None self._retry: Optional[int] = None self.max_buffer = max_buffer self.events_parsed = 0 self.bytes_fed = 0 self._bom_checked = False def feed(self, chunk: bytes) -> list[SSEEvent]: self.bytes_fed += len(chunk) self._buf += chunk if len(self._buf) > self.max_buffer: raise SSEBufferOverflow(f"single SSE event exceeds {self.max_buffer} bytes") out: list[SSEEvent] = [] while True: nl = self._buf.find(b"\n") if nl < 0: break raw = bytes(self._buf[:nl]) del self._buf[: nl + 1] if raw.endswith(b"\r"): raw = raw[:-1] if not self._bom_checked: self._bom_checked = True if raw.startswith(b"\xef\xbb\xbf"): raw = raw[3:] ev = self._line(raw.decode("utf-8", "replace")) if ev is not None: out.append(ev) return out def close(self) -> list[SSEEvent]: """Flush a trailing event that lacked its blank line (server closed the connection).""" out: list[SSEEvent] = [] if self._buf: ev = self._line(bytes(self._buf).decode("utf-8", "replace").rstrip("\r")) self._buf.clear() if ev: out.append(ev) if self._data: out.append(self._dispatch()) return out def _line(self, line: str) -> Optional[SSEEvent]: if line == "": return self._dispatch() if self._data or self._event else None if line.startswith(":"): return None # comment / keep-alive name, sep, value = line.partition(":") if sep and value.startswith(" "): value = value[1:] if name == "event": self._event = value elif name == "data": self._data.append(value) elif name == "id": if "\x00" not in value: self._id = value elif name == "retry": if value.isdigit(): self._retry = int(value) # unknown fields are ignored per spec return None def _dispatch(self) -> SSEEvent: ev = SSEEvent(self._event, "\n".join(self._data), self._id, self._retry) self._event, self._data = None, [] self.events_parsed += 1 return ev def parse_sse_bytes(data: bytes, chunk_size: Optional[int] = None) -> list[SSEEvent]: """Parse a whole recorded stream (optionally in fixed-size chunks to exercise boundary handling).""" p = SSEParser() out: list[SSEEvent] = [] if chunk_size: for i in range(0, len(data), chunk_size): out += p.feed(data[i:i + chunk_size]) else: out += p.feed(data) out += p.close() return out class JSONArrayStreamParser: """Incremental parser for Gemini's default (non-SSE) streaming body: `application/json`, a JSON array of GenerateContentResponse objects that is pretty-printed and flushed object by object (docs/gemini/streaming.md §1: line 1 `[{`, objects separated by a line containing only `,`, closed by `]`). Emits each complete top-level object as an SSEEvent(event=None, data=) so the same normalisers apply. Bounded like SSEParser.""" def __init__(self, max_buffer: int = 8 * 1024 * 1024): self._buf = bytearray() self._depth = 0 self._in_str = False self._esc = False self._start: Optional[int] = None # offset (in _buf) of the object currently being read self._scanned = 0 # bytes of _buf already scanned self.max_buffer = max_buffer self.events_parsed = 0 def feed(self, chunk: bytes) -> list[SSEEvent]: out: list[SSEEvent] = [] buf = self._buf buf += chunk if len(buf) > self.max_buffer: raise SSEBufferOverflow(f"single JSON object exceeds {self.max_buffer} bytes") i = self._scanned consumed = 0 if self._start is not None else self._scanned while i < len(buf): b = buf[i] if self._depth == 0 and not self._in_str: if b == 0x7B: # `{` — a top-level object starts self._start = i self._depth = 1 else: # `[`, `]`, `,`, whitespace between objects → discard consumed = i + 1 i += 1 continue if self._in_str: if self._esc: self._esc = False elif b == 0x5C: self._esc = True elif b == 0x22: self._in_str = False elif b == 0x22: self._in_str = True elif b == 0x7B: self._depth += 1 elif b == 0x7D: self._depth -= 1 if self._depth == 0 and self._start is not None: out.append(SSEEvent(None, bytes(buf[self._start:i + 1]).decode("utf-8", "replace"))) self.events_parsed += 1 consumed = i + 1 self._start = None i += 1 if self._start is not None: # partial object: keep from its start, everything kept is already scanned del buf[:self._start] self._start = 0 else: del buf[:consumed] self._scanned = len(buf) return out def close(self) -> list[SSEEvent]: self._buf.clear() return [] # -------------------------------------------------------------------------------------- # 2. Partial JSON accumulation (tool arguments) # -------------------------------------------------------------------------------------- def parse_partial_json(text: str) -> Any: """Best-effort parse of a JSON prefix: closes open strings/arrays/objects. Returns None when hopeless. Only for UI previews — the authoritative value is the final `arguments`/`input` once the block is done.""" if not text.strip(): return None try: return json.loads(text) except json.JSONDecodeError: pass stack: list[str] = [] in_str = esc = False for ch in text: if in_str: if esc: esc = False elif ch == "\\": esc = True elif ch == '"': in_str = False elif ch == '"': in_str = True elif ch in "{[": stack.append("}" if ch == "{" else "]") elif ch in "}]" and stack: stack.pop() fixed = text if in_str: fixed += '"' fixed = fixed.rstrip() def _drop_dangling_key(s: str) -> str: """Remove a trailing `"key"` (with optional `:`) that has no value yet, plus any trailing comma.""" s = s.rstrip().rstrip(",").rstrip() if s.endswith(":"): s = s[:-1].rstrip() if s.endswith('"'): j = s.rfind('"', 0, len(s) - 1) before = s[:j].rstrip() if before.endswith(("{", ",")): # the string is an object key, not a value s = before.rstrip(",").rstrip() return s for _ in range(3): try: return json.loads(fixed + "".join(reversed(stack))) except json.JSONDecodeError: nxt = _drop_dangling_key(fixed) if nxt == fixed: break fixed = nxt return None @dataclass class PartialJSON: """Per-tool-call argument accumulator.""" buffers: dict[str, str] = field(default_factory=dict) def append(self, key: str, fragment: str) -> str: self.buffers[key] = self.buffers.get(key, "") + fragment return self.buffers[key] def preview(self, key: str) -> Any: return parse_partial_json(self.buffers.get(key, "")) def finish(self, key: str, final: Optional[str] = None) -> dict[str, Any]: raw = final if final is not None else self.buffers.pop(key, "") self.buffers.pop(key, None) if not raw.strip(): return {} try: v = json.loads(raw) return v if isinstance(v, dict) else {"_value": v} except json.JSONDecodeError: return {"_raw": raw, "_error": "invalid JSON"} # -------------------------------------------------------------------------------------- # 3. Unified normalisers # -------------------------------------------------------------------------------------- @dataclass class Unified: type: str # text_delta | reasoning_delta | tool_call_start | tool_call_delta | tool_call_done | usage | done | error | raw text: str = "" tool_call_id: Optional[str] = None tool_name: Optional[str] = None arguments_delta: str = "" arguments: Optional[dict[str, Any]] = None usage: Optional[dict[str, Any]] = None stop_reason: Optional[str] = None error: Optional[dict[str, Any]] = None sequence: Optional[int] = None raw: Optional[dict[str, Any]] = None class OpenAINormalizer: """Responses API events → Unified. Tracks function_call items by item_id; `sequence` = sequence_number (resume cursor).""" TERMINAL = {"response.completed": "end", "response.incomplete": "incomplete", "response.failed": "failed", "response.cancelled": "cancelled"} def __init__(self) -> None: self.args = PartialJSON() self.items: dict[str, dict[str, Any]] = {} self.last_sequence: Optional[int] = None self.response_id: Optional[str] = None def __call__(self, ev: SSEEvent) -> list[Unified]: obj = ev.json() if not isinstance(obj, dict): return [] if ev.data == "[DONE]" else [Unified("raw", raw={"event": ev.event, "data": ev.data})] t = obj.get("type") or ev.event or "" seq = obj.get("sequence_number") if isinstance(seq, int): self.last_sequence = seq if self.response_id is None and isinstance(obj.get("response"), dict): self.response_id = obj["response"].get("id") u = lambda typ, **kw: [Unified(typ, sequence=seq, raw=obj, **kw)] # noqa: E731 if t == "response.output_text.delta": return u("text_delta", text=obj.get("delta", "")) if t in ("response.reasoning_summary_text.delta", "response.reasoning_text.delta"): return u("reasoning_delta", text=obj.get("delta", "")) if t == "response.refusal.delta": return u("text_delta", text=obj.get("delta", "")) if t == "response.output_item.added": item = obj.get("item") or {} if item.get("type") in ("function_call", "custom_tool_call", "mcp_call"): self.items[item.get("id")] = item return u("tool_call_start", tool_call_id=item.get("call_id") or item.get("id"), tool_name=item.get("name")) return u("raw") if t in ("response.function_call_arguments.delta", "response.custom_tool_call_input.delta", "response.mcp_call_arguments.delta"): iid = obj.get("item_id") self.args.append(iid, obj.get("delta", "")) return u("tool_call_delta", tool_call_id=self._call_id(iid), tool_name=(self.items.get(iid) or {}).get("name"), arguments_delta=obj.get("delta", "")) if t in ("response.function_call_arguments.done", "response.custom_tool_call_input.done", "response.mcp_call_arguments.done"): iid = obj.get("item_id") final = obj.get("arguments") if "arguments" in obj else obj.get("input") args = self.args.finish(iid, final if isinstance(final, str) else None) return u("tool_call_done", tool_call_id=self._call_id(iid), tool_name=(self.items.get(iid) or {}).get("name"), arguments=args) if t in self.TERMINAL: resp = obj.get("response") or {} out = [] if isinstance(resp.get("usage"), dict): out += u("usage", usage=resp["usage"]) stop = self.TERMINAL[t] if t == "response.incomplete" and (resp.get("incomplete_details") or {}).get("reason") == "max_output_tokens": stop = "max_tokens" return out + u("done", stop_reason=stop) if t == "error": return u("error", error=obj) return u("raw") def _call_id(self, item_id: Optional[str]) -> Optional[str]: item = self.items.get(item_id) or {} return item.get("call_id") or item_id class AnthropicNormalizer: """Messages API events → Unified. Tracks content blocks by index; input_json_delta fragments are accumulated.""" STOP = {"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"} MESSAGE_EVENTS = {"message_start", "content_block_start", "content_block_delta", "content_block_stop", "message_delta", "message_stop", "ping"} def __init__(self) -> None: self.args = PartialJSON() self.blocks: dict[int, dict[str, Any]] = {} self.usage: dict[str, Any] = {} self.stop_reason: Optional[str] = None self.message_id: Optional[str] = None def __call__(self, ev: SSEEvent) -> list[Unified]: obj = ev.json() if not isinstance(obj, dict): return [Unified("raw", raw={"event": ev.event, "data": ev.data})] t = obj.get("type") or ev.event or "" u = lambda typ, **kw: [Unified(typ, raw=obj, **kw)] # noqa: E731 if t == "message_start": msg = obj.get("message") or {} self.message_id = msg.get("id") self.usage = dict(msg.get("usage") or {}) return u("raw") if t == "content_block_start": idx, cb = obj.get("index"), obj.get("content_block") or {} self.blocks[idx] = cb if cb.get("type") in ("tool_use", "server_tool_use", "mcp_tool_use"): return u("tool_call_start", tool_call_id=cb.get("id"), tool_name=cb.get("name")) return u("raw") if t == "content_block_delta": idx, d = obj.get("index"), obj.get("delta") or {} dt = d.get("type") if dt == "text_delta": return u("text_delta", text=d.get("text", "")) if dt == "thinking_delta": return u("reasoning_delta", text=d.get("thinking", "")) if dt == "input_json_delta": cb = self.blocks.get(idx) or {} frag = d.get("partial_json", "") self.args.append(str(idx), frag) return u("tool_call_delta", tool_call_id=cb.get("id"), tool_name=cb.get("name"), arguments_delta=frag) return u("raw") # signature_delta, citations_delta if t == "content_block_stop": idx = obj.get("index") cb = self.blocks.get(idx) or {} if cb.get("type") in ("tool_use", "server_tool_use", "mcp_tool_use"): args = self.args.finish(str(idx)) if not args and isinstance(cb.get("input"), dict) and cb["input"]: args = cb["input"] # some server tools send the full input in content_block_start return u("tool_call_done", tool_call_id=cb.get("id"), tool_name=cb.get("name"), arguments=args) return u("raw") if t == "message_delta": self.stop_reason = self.STOP.get((obj.get("delta") or {}).get("stop_reason") or "", None) self.usage.update(obj.get("usage") or {}) return u("usage", usage=dict(self.usage)) if t == "message_stop": return u("done", stop_reason=self.stop_reason or "end") if t == "error": return u("error", error=obj.get("error") or obj) return u("raw") # ping and future event types class XAINormalizer: """xAI streams in three framings; dispatch on the payload shape: * Chat Completions: `data:`-only chunks `{"object":"chat.completion.chunk","choices":[{"delta":{…},"finish_reason"}]}` — `delta.reasoning_content` (BEFORE content) → reasoning_delta, `delta.content` → text_delta, `delta.tool_calls[]` arrive WHOLE in one chunk (start+delta+done emitted together), final `choices:[]` chunk carries `usage` when `stream_options.include_usage` was set, then `data: [DONE]` → done (stop reason from the last finish_reason). * Responses API: OpenAI event vocabulary with `sequence_number` — delegated to OpenAINormalizer (incl. `response.reasoning_summary_text.delta`, code_interpreter events surfaced as raw). * /v1/messages (Anthropic-compatible, DEPRECATED): delegated to AnthropicNormalizer; note live quirks — `index` is 0 for both the thinking and the text block and deltas carry no `index` (the delegate tolerates both).""" CHAT_FINISH = {"stop": "end", "length": "max_tokens", "tool_calls": "tool_use", "content_filter": "refusal"} def __init__(self) -> None: self.responses = OpenAINormalizer() self.messages = AnthropicNormalizer() self.chat_id: Optional[str] = None self.chat_stop: Optional[str] = None self.saw_chat = False def __call__(self, ev: SSEEvent) -> list[Unified]: obj = ev.json() if not isinstance(obj, dict): if ev.data == "[DONE]": return [Unified("done", stop_reason=self.chat_stop or "end", raw={"data": "[DONE]"})] if self.saw_chat else [] return [Unified("raw", raw={"event": ev.event, "data": ev.data})] t = obj.get("type") or ev.event or "" if obj.get("object") == "chat.completion.chunk" or ("choices" in obj and not t): return self._chat(obj) if t.startswith("response.") or t in ("error",) and "sequence_number" in obj: return self.responses(ev) if t in AnthropicNormalizer.MESSAGE_EVENTS: return self.messages(ev) if t == "error": # bare error object ({"code","error"} shape) — xAI returns HTTP errors before any chunk, but be safe return [Unified("error", error=obj, raw=obj)] return [Unified("raw", raw=obj)] def _chat(self, obj: dict[str, Any]) -> list[Unified]: self.saw_chat = True self.chat_id = obj.get("id") or self.chat_id out: list[Unified] = [] choices = obj.get("choices") or [] if not choices and isinstance(obj.get("usage"), dict): return [Unified("usage", usage=obj["usage"], raw=obj)] for ch in choices: d = ch.get("delta") or {} if d.get("reasoning_content"): out.append(Unified("reasoning_delta", text=d["reasoning_content"], raw=obj)) if d.get("content"): out.append(Unified("text_delta", text=d["content"], raw=obj)) for tc in d.get("tool_calls") or []: fn = tc.get("function") or {} raw_args = fn.get("arguments") or "" cid = tc.get("id") or f"{self.chat_id}:{tc.get('index', 0)}" out.append(Unified("tool_call_start", tool_call_id=cid, tool_name=fn.get("name"), raw=obj)) if raw_args: out.append(Unified("tool_call_delta", tool_call_id=cid, tool_name=fn.get("name"), arguments_delta=raw_args, raw=obj)) args = PartialJSON().finish(cid, raw_args or None) if raw_args else {} out.append(Unified("tool_call_done", tool_call_id=cid, tool_name=fn.get("name"), arguments=args, raw=obj)) if ch.get("finish_reason"): self.chat_stop = self.CHAT_FINISH.get(ch["finish_reason"], ch["finish_reason"]) if not out: out.append(Unified("raw", raw=obj)) return out class GeminiNormalizer: """Gemini streams → Unified. Two payload families share one normaliser: * generateContent chunks (SSE `?alt=sse` or JSON-array framing): `candidates[0].content.parts[]` deltas — `{text}` → text_delta, `{text, thought:true}` → reasoning_delta, `{functionCall{name,args,id}}` (+ sibling `thoughtSignature`) → tool_call_start/delta/done at once (arguments never arrive partially), empty-text part with `thoughtSignature` → raw (keep it: it must be echoed in multi-turn), `inlineData`/`executableCode`/`codeExecutionResult` → raw. `usageMetadata` is on EVERY chunk (running totals; the last wins) → one `usage` event when `finishReason` arrives, then `done` (STOP→end / tool_use if functionCall seen, MAX_TOKENS→max_tokens, SAFETY & co→refusal). A chunk with `promptFeedback` and no candidates = prompt blocked → done(refusal). `{"error": google.rpc.Status}` → error. * Interactions API events (`event_type`): interaction.created/status_update → raw; step.start(function_call) → tool_call_start; step.delta text → text_delta, thought_summary → reasoning_delta, thought_signature → raw, arguments_delta → tool_call_delta; step.stop of a function_call step → tool_call_done; interaction.completed → usage + done (status completed→end, requires_action→tool_use, incomplete→incomplete, failed→failed, cancelled→cancelled); `error` → error; `event: done` / `data: [DONE]` → nothing (terminal marker after interaction.completed). `last_signature` keeps the most recent thoughtSignature for callers that rebuild the model turn by hand.""" REFUSAL = {"SAFETY", "RECITATION", "BLOCKLIST", "PROHIBITED_CONTENT", "SPII", "IMAGE_SAFETY", "IMAGE_PROHIBITED_CONTENT", "IMAGE_RECITATION", "IMAGE_OTHER", "LANGUAGE"} INTERACTION_STATUS = {"completed": "end", "requires_action": "tool_use", "incomplete": "incomplete", "failed": "failed", "cancelled": "cancelled"} def __init__(self) -> None: self.args = PartialJSON() self.usage: dict[str, Any] = {} self.finish_reason: Optional[str] = None self.saw_call = False self.response_id: Optional[str] = None self.model_version: Optional[str] = None self.last_signature: Optional[str] = None self.steps: dict[int, dict[str, Any]] = {} # Interactions: step index → step header self.interaction_id: Optional[str] = None def __call__(self, ev: SSEEvent) -> list[Unified]: obj = ev.json() if not isinstance(obj, dict): return [] if ev.data == "[DONE]" else [Unified("raw", raw={"event": ev.event, "data": ev.data})] if "event_type" in obj or ev.event in ("step.delta", "step.start", "step.stop", "interaction.created", "interaction.completed"): return self._interaction(obj, obj.get("event_type") or ev.event or "") return self._chunk(obj) def _stop(self) -> str: fr = self.finish_reason if self.saw_call: return "tool_use" if fr in (None, "STOP"): return "end" if fr == "MAX_TOKENS": return "max_tokens" if fr in self.REFUSAL or fr == "BLOCKED": return "refusal" return "other" # MALFORMED_FUNCTION_CALL, MISSING_THOUGHT_SIGNATURE, UNEXPECTED_TOOL_CALL, TOO_MANY_TOOL_CALLS, MALFORMED_RESPONSE… def _chunk(self, obj: dict[str, Any]) -> list[Unified]: u = lambda typ, **kw: Unified(typ, raw=obj, **kw) # noqa: E731 if isinstance(obj.get("error"), dict): return [u("error", error=obj["error"])] if isinstance(obj.get("usageMetadata"), dict): self.usage = dict(obj["usageMetadata"]) self.response_id = obj.get("responseId") or self.response_id self.model_version = obj.get("modelVersion") or self.model_version cands = obj.get("candidates") or [] out: list[Unified] = [] if not cands: if obj.get("promptFeedback"): self.finish_reason = "BLOCKED" out.append(u("usage", usage=dict(self.usage))) out.append(u("done", stop_reason="refusal")) return out return [u("raw")] c = cands[0] for i, p in enumerate((c.get("content") or {}).get("parts") or []): if p.get("thoughtSignature"): self.last_signature = p["thoughtSignature"] if "functionCall" in p: fc = p["functionCall"] or {} args = fc.get("args") or {} cid = fc.get("id") or f"call_{len(self.steps)}_{i}" self.saw_call = True out.append(u("tool_call_start", tool_call_id=cid, tool_name=fc.get("name"))) out.append(u("tool_call_delta", tool_call_id=cid, tool_name=fc.get("name"), arguments_delta=json.dumps(args))) out.append(u("tool_call_done", tool_call_id=cid, tool_name=fc.get("name"), arguments=args if isinstance(args, dict) else {"_value": args})) elif "text" in p: if p.get("thought"): out.append(u("reasoning_delta", text=p.get("text") or "")) elif p.get("text"): out.append(u("text_delta", text=p["text"])) else: out.append(u("raw")) # empty-text signature carrier (Gemini 3 last chunk) else: out.append(u("raw")) # inlineData, executableCode, codeExecutionResult, fileData… if c.get("finishReason"): self.finish_reason = c["finishReason"] out.append(u("usage", usage=dict(self.usage))) out.append(u("done", stop_reason=self._stop())) return out or [u("raw")] def _interaction(self, obj: dict[str, Any], t: str) -> list[Unified]: u = lambda typ, **kw: [Unified(typ, raw=obj, **kw)] # noqa: E731 idx = obj.get("index") if t == "interaction.created": self.interaction_id = (obj.get("interaction") or {}).get("id") return u("raw") if t == "step.start": step = obj.get("step") or {} self.steps[idx] = step if step.get("type") == "function_call": return u("tool_call_start", tool_call_id=step.get("id"), tool_name=step.get("name")) return u("raw") if t == "step.delta": d = obj.get("delta") or {} dt = d.get("type") step = self.steps.get(idx) or {} if dt == "text": return u("text_delta", text=d.get("text") or "") if dt == "thought_summary": return u("reasoning_delta", text=((d.get("content") or {}).get("text") or "")) if dt == "thought_signature": self.last_signature = d.get("signature") or self.last_signature return u("raw") if dt == "arguments_delta": frag = d.get("arguments") or "" self.args.append(str(idx), frag) return u("tool_call_delta", tool_call_id=step.get("id"), tool_name=step.get("name"), arguments_delta=frag) return u("raw") # image/audio/video/document/function_result/*_call/*_result/processing_*/text_annotation_delta if t == "step.stop": step = self.steps.get(idx) or {} if step.get("type") == "function_call": args = self.args.finish(str(idx)) if not args and isinstance(step.get("arguments"), dict) and step["arguments"]: args = step["arguments"] return u("tool_call_done", tool_call_id=step.get("id"), tool_name=step.get("name"), arguments=args) return u("raw") if t == "interaction.completed": inter = obj.get("interaction") or {} out = [] if isinstance(inter.get("usage"), dict): out += u("usage", usage=inter["usage"]) return out + u("done", stop_reason=self.INTERACTION_STATUS.get(inter.get("status") or "", inter.get("status") or "end")) if t == "error": return u("error", error=obj.get("error") or obj) return u("raw") # interaction.status_update / in_progress / requires_action and future types def normalizer_for(provider: str): return {"openai": OpenAINormalizer, "anthropic": AnthropicNormalizer, "xai": XAINormalizer, "gemini": GeminiNormalizer}[provider]() def iter_unified(chunks: Iterable[bytes], provider: str, framing: str = "sse") -> Iterator[Unified]: """chunks: iterable of raw byte chunks from the HTTP body (any chunking). framing: "sse" (all providers) or "json_array" (Gemini streamGenerateContent WITHOUT `?alt=sse`).""" parser = JSONArrayStreamParser() if framing == "json_array" else SSEParser() norm = normalizer_for(provider) for chunk in chunks: for ev in parser.feed(chunk): yield from norm(ev) for ev in parser.close(): yield from norm(ev) def collect_text(events: Iterable[Unified]) -> str: return "".join(e.text for e in events if e.type == "text_delta") if __name__ == "__main__": import sys if len(sys.argv) not in (3, 4): print("usage: sse_parser.py [sse|json_array]") sys.exit(1) data = open(sys.argv[2], "rb").read() for u in iter_unified([data[i:i + 64] for i in range(0, len(data), 64)], sys.argv[1], sys.argv[3] if len(sys.argv) == 4 else "sse"): if u.type != "raw": print(u.type, u.text or u.arguments_delta or u.arguments or u.usage or u.stop_reason or u.error or "")