Python 88.3%
TypeScript 7.6%
Shell 4.1%
1#!/usr/bin/env python32"""Generic Server-Sent-Events parser + provider normalisers for OpenAI (Responses API), Anthropic (Messages API),3xAI (Chat Completions chunks · Responses events · Anthropic-compatible /v1/messages) and Google Gemini4(streamGenerateContent `?alt=sse` · JSON-array framing · Interactions API `step.*` events).56STATUS: DOCUMENTED (framing rules from the WHATWG SSE spec, event vocabularies from the OpenAPI specs / discovery document /7build-with-claude/streaming.md) · offline-tested on recorded fixtures (tests/shared/fixtures/sse/*) by8tests/shared/test_sse_parser.py (2026-09-18 OpenAI/Anthropic; 2026-09-19 xAI/Gemini fixtures cut from live captures in9tmp-live/xai/*_stream*.json and tmp-live/gemini-core/*.raw.txt, thoughtSignature blobs truncated). Stdlib only.10See docs/architecture/streaming-patterns.md.1112Layers13 1. SSEParser — incremental, byte-oriented, bounded buffer (backpressure friendly): feed(bytes) -> list[SSEEvent]14 Handles `event:`/`data:`/`id:`/`retry:` fields, comments (`: keep-alive`), multi-line data,15 CRLF/LF, BOM, chunk boundaries splitting a line, and `data: [DONE]` (Chat Completions sentinel).16 JSONArrayStreamParser — Gemini's NON-SSE streaming mode (no `?alt=sse`): a pretty-printed JSON array streamed17 incrementally (`[{` … `}\n,\n{` … `}\n]`); emits each top-level object as soon as it closes.18 2. PartialJSON — accumulates argument fragments per tool call and offers a tolerant preview parse.19 3. OpenAINormalizer / AnthropicNormalizer / XAINormalizer / GeminiNormalizer — unified vocabulary:20 text_delta · reasoning_delta · tool_call_start · tool_call_delta · tool_call_done · usage · done · error · raw21 4. iter_unified(chunks, provider, framing="sse"|"json_array") — convenience generator over an iterable of byte chunks.2223Framing differences (the JSON payload is authoritative; `event:` names are absent on two of the four providers):24 OpenAI event: response.output_text.delta data: {"type":"response.output_text.delta","sequence_number":7,"item_id":…,"delta":"Hi"}25 Anthropic event: content_block_delta data: {"type":"content_block_delta","index":0,"delta":{"type":"text_delta","text":"Hi"}}26 xAI chat data: {"object":"chat.completion.chunk","choices":[{"delta":{"reasoning_content":"…"}}]} … data: [DONE] (no event: lines)27 xAI resp. event: response.output_text.delta (identical to OpenAI incl. sequence_number; `response.reasoning_summary_text.delta` observed)28 Gemini data: {"candidates":[{"content":{"parts":[{"text":"Hi"}]}}],"usageMetadata":{…}}\r\n\r\n (no event:, no [DONE];29 the LAST chunk has finishReason and, on Gemini 3, an EMPTY text part carrying `thoughtSignature`)30 Gemini Interactions event: step.delta data: {"event_type":"step.delta","index":1,"delta":{"type":"text","text":"OK."}} … event: done / data: [DONE]31"""32from __future__ import annotations3334import json35from dataclasses import dataclass, field36from typing import Any, Iterable, Iterator, Optional373839# --------------------------------------------------------------------------------------40# 1. Generic SSE parser41# --------------------------------------------------------------------------------------4243@dataclass44class SSEEvent:45 event: Optional[str]46 data: str47 id: Optional[str] = None48 retry: Optional[int] = None4950 def json(self) -> Any:51 try:52 return json.loads(self.data)53 except json.JSONDecodeError:54 return None555657class SSEBufferOverflow(Exception):58 pass596061class SSEParser:62 """Incremental SSE parser. Call feed() with each network chunk (any size); events are returned as soon as their63 terminating blank line arrives. `max_buffer` bounds memory for a single (possibly hostile) event."""6465 def __init__(self, max_buffer: int = 8 * 1024 * 1024):66 self._buf = bytearray()67 self._event: Optional[str] = None68 self._data: list[str] = []69 self._id: Optional[str] = None70 self._retry: Optional[int] = None71 self.max_buffer = max_buffer72 self.events_parsed = 073 self.bytes_fed = 074 self._bom_checked = False7576 def feed(self, chunk: bytes) -> list[SSEEvent]:77 self.bytes_fed += len(chunk)78 self._buf += chunk79 if len(self._buf) > self.max_buffer:80 raise SSEBufferOverflow(f"single SSE event exceeds {self.max_buffer} bytes")81 out: list[SSEEvent] = []82 while True:83 nl = self._buf.find(b"\n")84 if nl < 0:85 break86 raw = bytes(self._buf[:nl])87 del self._buf[: nl + 1]88 if raw.endswith(b"\r"):89 raw = raw[:-1]90 if not self._bom_checked:91 self._bom_checked = True92 if raw.startswith(b"\xef\xbb\xbf"):93 raw = raw[3:]94 ev = self._line(raw.decode("utf-8", "replace"))95 if ev is not None:96 out.append(ev)97 return out9899 def close(self) -> list[SSEEvent]:100 """Flush a trailing event that lacked its blank line (server closed the connection)."""101 out: list[SSEEvent] = []102 if self._buf:103 ev = self._line(bytes(self._buf).decode("utf-8", "replace").rstrip("\r"))104 self._buf.clear()105 if ev:106 out.append(ev)107 if self._data:108 out.append(self._dispatch())109 return out110111 def _line(self, line: str) -> Optional[SSEEvent]:112 if line == "":113 return self._dispatch() if self._data or self._event else None114 if line.startswith(":"):115 return None # comment / keep-alive116 name, sep, value = line.partition(":")117 if sep and value.startswith(" "):118 value = value[1:]119 if name == "event":120 self._event = value121 elif name == "data":122 self._data.append(value)123 elif name == "id":124 if "\x00" not in value:125 self._id = value126 elif name == "retry":127 if value.isdigit():128 self._retry = int(value)129 # unknown fields are ignored per spec130 return None131132 def _dispatch(self) -> SSEEvent:133 ev = SSEEvent(self._event, "\n".join(self._data), self._id, self._retry)134 self._event, self._data = None, []135 self.events_parsed += 1136 return ev137138139def parse_sse_bytes(data: bytes, chunk_size: Optional[int] = None) -> list[SSEEvent]:140 """Parse a whole recorded stream (optionally in fixed-size chunks to exercise boundary handling)."""141 p = SSEParser()142 out: list[SSEEvent] = []143 if chunk_size:144 for i in range(0, len(data), chunk_size):145 out += p.feed(data[i:i + chunk_size])146 else:147 out += p.feed(data)148 out += p.close()149 return out150151152class JSONArrayStreamParser:153 """Incremental parser for Gemini's default (non-SSE) streaming body: `application/json`, a JSON array of154 GenerateContentResponse objects that is pretty-printed and flushed object by object (docs/gemini/streaming.md §1:155 line 1 `[{`, objects separated by a line containing only `,`, closed by `]`). Emits each complete top-level object as an156 SSEEvent(event=None, data=<object JSON>) so the same normalisers apply. Bounded like SSEParser."""157158 def __init__(self, max_buffer: int = 8 * 1024 * 1024):159 self._buf = bytearray()160 self._depth = 0161 self._in_str = False162 self._esc = False163 self._start: Optional[int] = None # offset (in _buf) of the object currently being read164 self._scanned = 0 # bytes of _buf already scanned165 self.max_buffer = max_buffer166 self.events_parsed = 0167168 def feed(self, chunk: bytes) -> list[SSEEvent]:169 out: list[SSEEvent] = []170 buf = self._buf171 buf += chunk172 if len(buf) > self.max_buffer:173 raise SSEBufferOverflow(f"single JSON object exceeds {self.max_buffer} bytes")174 i = self._scanned175 consumed = 0 if self._start is not None else self._scanned176 while i < len(buf):177 b = buf[i]178 if self._depth == 0 and not self._in_str:179 if b == 0x7B: # `{` — a top-level object starts180 self._start = i181 self._depth = 1182 else: # `[`, `]`, `,`, whitespace between objects → discard183 consumed = i + 1184 i += 1185 continue186 if self._in_str:187 if self._esc:188 self._esc = False189 elif b == 0x5C:190 self._esc = True191 elif b == 0x22:192 self._in_str = False193 elif b == 0x22:194 self._in_str = True195 elif b == 0x7B:196 self._depth += 1197 elif b == 0x7D:198 self._depth -= 1199 if self._depth == 0 and self._start is not None:200 out.append(SSEEvent(None, bytes(buf[self._start:i + 1]).decode("utf-8", "replace")))201 self.events_parsed += 1202 consumed = i + 1203 self._start = None204 i += 1205 if self._start is not None: # partial object: keep from its start, everything kept is already scanned206 del buf[:self._start]207 self._start = 0208 else:209 del buf[:consumed]210 self._scanned = len(buf)211 return out212213 def close(self) -> list[SSEEvent]:214 self._buf.clear()215 return []216217218# --------------------------------------------------------------------------------------219# 2. Partial JSON accumulation (tool arguments)220# --------------------------------------------------------------------------------------221222def parse_partial_json(text: str) -> Any:223 """Best-effort parse of a JSON prefix: closes open strings/arrays/objects. Returns None when hopeless.224 Only for UI previews — the authoritative value is the final `arguments`/`input` once the block is done."""225 if not text.strip():226 return None227 try:228 return json.loads(text)229 except json.JSONDecodeError:230 pass231 stack: list[str] = []232 in_str = esc = False233 for ch in text:234 if in_str:235 if esc:236 esc = False237 elif ch == "\\":238 esc = True239 elif ch == '"':240 in_str = False241 elif ch == '"':242 in_str = True243 elif ch in "{[":244 stack.append("}" if ch == "{" else "]")245 elif ch in "}]" and stack:246 stack.pop()247 fixed = text248 if in_str:249 fixed += '"'250 fixed = fixed.rstrip()251252 def _drop_dangling_key(s: str) -> str:253 """Remove a trailing `"key"` (with optional `:`) that has no value yet, plus any trailing comma."""254 s = s.rstrip().rstrip(",").rstrip()255 if s.endswith(":"):256 s = s[:-1].rstrip()257 if s.endswith('"'):258 j = s.rfind('"', 0, len(s) - 1)259 before = s[:j].rstrip()260 if before.endswith(("{", ",")): # the string is an object key, not a value261 s = before.rstrip(",").rstrip()262 return s263264 for _ in range(3):265 try:266 return json.loads(fixed + "".join(reversed(stack)))267 except json.JSONDecodeError:268 nxt = _drop_dangling_key(fixed)269 if nxt == fixed:270 break271 fixed = nxt272 return None273274275@dataclass276class PartialJSON:277 """Per-tool-call argument accumulator."""278 buffers: dict[str, str] = field(default_factory=dict)279280 def append(self, key: str, fragment: str) -> str:281 self.buffers[key] = self.buffers.get(key, "") + fragment282 return self.buffers[key]283284 def preview(self, key: str) -> Any:285 return parse_partial_json(self.buffers.get(key, ""))286287 def finish(self, key: str, final: Optional[str] = None) -> dict[str, Any]:288 raw = final if final is not None else self.buffers.pop(key, "")289 self.buffers.pop(key, None)290 if not raw.strip():291 return {}292 try:293 v = json.loads(raw)294 return v if isinstance(v, dict) else {"_value": v}295 except json.JSONDecodeError:296 return {"_raw": raw, "_error": "invalid JSON"}297298299# --------------------------------------------------------------------------------------300# 3. Unified normalisers301# --------------------------------------------------------------------------------------302303@dataclass304class Unified:305 type: str # text_delta | reasoning_delta | tool_call_start | tool_call_delta | tool_call_done | usage | done | error | raw306 text: str = ""307 tool_call_id: Optional[str] = None308 tool_name: Optional[str] = None309 arguments_delta: str = ""310 arguments: Optional[dict[str, Any]] = None311 usage: Optional[dict[str, Any]] = None312 stop_reason: Optional[str] = None313 error: Optional[dict[str, Any]] = None314 sequence: Optional[int] = None315 raw: Optional[dict[str, Any]] = None316317318class OpenAINormalizer:319 """Responses API events → Unified. Tracks function_call items by item_id; `sequence` = sequence_number (resume cursor)."""320321 TERMINAL = {"response.completed": "end", "response.incomplete": "incomplete", "response.failed": "failed", "response.cancelled": "cancelled"}322323 def __init__(self) -> None:324 self.args = PartialJSON()325 self.items: dict[str, dict[str, Any]] = {}326 self.last_sequence: Optional[int] = None327 self.response_id: Optional[str] = None328329 def __call__(self, ev: SSEEvent) -> list[Unified]:330 obj = ev.json()331 if not isinstance(obj, dict):332 return [] if ev.data == "[DONE]" else [Unified("raw", raw={"event": ev.event, "data": ev.data})]333 t = obj.get("type") or ev.event or ""334 seq = obj.get("sequence_number")335 if isinstance(seq, int):336 self.last_sequence = seq337 if self.response_id is None and isinstance(obj.get("response"), dict):338 self.response_id = obj["response"].get("id")339 u = lambda typ, **kw: [Unified(typ, sequence=seq, raw=obj, **kw)] # noqa: E731340 if t == "response.output_text.delta":341 return u("text_delta", text=obj.get("delta", ""))342 if t in ("response.reasoning_summary_text.delta", "response.reasoning_text.delta"):343 return u("reasoning_delta", text=obj.get("delta", ""))344 if t == "response.refusal.delta":345 return u("text_delta", text=obj.get("delta", ""))346 if t == "response.output_item.added":347 item = obj.get("item") or {}348 if item.get("type") in ("function_call", "custom_tool_call", "mcp_call"):349 self.items[item.get("id")] = item350 return u("tool_call_start", tool_call_id=item.get("call_id") or item.get("id"), tool_name=item.get("name"))351 return u("raw")352 if t in ("response.function_call_arguments.delta", "response.custom_tool_call_input.delta", "response.mcp_call_arguments.delta"):353 iid = obj.get("item_id")354 self.args.append(iid, obj.get("delta", ""))355 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", ""))356 if t in ("response.function_call_arguments.done", "response.custom_tool_call_input.done", "response.mcp_call_arguments.done"):357 iid = obj.get("item_id")358 final = obj.get("arguments") if "arguments" in obj else obj.get("input")359 args = self.args.finish(iid, final if isinstance(final, str) else None)360 return u("tool_call_done", tool_call_id=self._call_id(iid), tool_name=(self.items.get(iid) or {}).get("name"), arguments=args)361 if t in self.TERMINAL:362 resp = obj.get("response") or {}363 out = []364 if isinstance(resp.get("usage"), dict):365 out += u("usage", usage=resp["usage"])366 stop = self.TERMINAL[t]367 if t == "response.incomplete" and (resp.get("incomplete_details") or {}).get("reason") == "max_output_tokens":368 stop = "max_tokens"369 return out + u("done", stop_reason=stop)370 if t == "error":371 return u("error", error=obj)372 return u("raw")373374 def _call_id(self, item_id: Optional[str]) -> Optional[str]:375 item = self.items.get(item_id) or {}376 return item.get("call_id") or item_id377378379class AnthropicNormalizer:380 """Messages API events → Unified. Tracks content blocks by index; input_json_delta fragments are accumulated."""381382 STOP = {"end_turn": "end", "max_tokens": "max_tokens", "tool_use": "tool_use", "stop_sequence": "stop_sequence",383 "refusal": "refusal", "pause_turn": "incomplete", "model_context_window_exceeded": "incomplete"}384 MESSAGE_EVENTS = {"message_start", "content_block_start", "content_block_delta", "content_block_stop", "message_delta", "message_stop", "ping"}385386 def __init__(self) -> None:387 self.args = PartialJSON()388 self.blocks: dict[int, dict[str, Any]] = {}389 self.usage: dict[str, Any] = {}390 self.stop_reason: Optional[str] = None391 self.message_id: Optional[str] = None392393 def __call__(self, ev: SSEEvent) -> list[Unified]:394 obj = ev.json()395 if not isinstance(obj, dict):396 return [Unified("raw", raw={"event": ev.event, "data": ev.data})]397 t = obj.get("type") or ev.event or ""398 u = lambda typ, **kw: [Unified(typ, raw=obj, **kw)] # noqa: E731399 if t == "message_start":400 msg = obj.get("message") or {}401 self.message_id = msg.get("id")402 self.usage = dict(msg.get("usage") or {})403 return u("raw")404 if t == "content_block_start":405 idx, cb = obj.get("index"), obj.get("content_block") or {}406 self.blocks[idx] = cb407 if cb.get("type") in ("tool_use", "server_tool_use", "mcp_tool_use"):408 return u("tool_call_start", tool_call_id=cb.get("id"), tool_name=cb.get("name"))409 return u("raw")410 if t == "content_block_delta":411 idx, d = obj.get("index"), obj.get("delta") or {}412 dt = d.get("type")413 if dt == "text_delta":414 return u("text_delta", text=d.get("text", ""))415 if dt == "thinking_delta":416 return u("reasoning_delta", text=d.get("thinking", ""))417 if dt == "input_json_delta":418 cb = self.blocks.get(idx) or {}419 frag = d.get("partial_json", "")420 self.args.append(str(idx), frag)421 return u("tool_call_delta", tool_call_id=cb.get("id"), tool_name=cb.get("name"), arguments_delta=frag)422 return u("raw") # signature_delta, citations_delta423 if t == "content_block_stop":424 idx = obj.get("index")425 cb = self.blocks.get(idx) or {}426 if cb.get("type") in ("tool_use", "server_tool_use", "mcp_tool_use"):427 args = self.args.finish(str(idx))428 if not args and isinstance(cb.get("input"), dict) and cb["input"]:429 args = cb["input"] # some server tools send the full input in content_block_start430 return u("tool_call_done", tool_call_id=cb.get("id"), tool_name=cb.get("name"), arguments=args)431 return u("raw")432 if t == "message_delta":433 self.stop_reason = self.STOP.get((obj.get("delta") or {}).get("stop_reason") or "", None)434 self.usage.update(obj.get("usage") or {})435 return u("usage", usage=dict(self.usage))436 if t == "message_stop":437 return u("done", stop_reason=self.stop_reason or "end")438 if t == "error":439 return u("error", error=obj.get("error") or obj)440 return u("raw") # ping and future event types441442443class XAINormalizer:444 """xAI streams in three framings; dispatch on the payload shape:445 * Chat Completions: `data:`-only chunks `{"object":"chat.completion.chunk","choices":[{"delta":{…},"finish_reason"}]}`446 — `delta.reasoning_content` (BEFORE content) → reasoning_delta, `delta.content` → text_delta, `delta.tool_calls[]` arrive447 WHOLE in one chunk (start+delta+done emitted together), final `choices:[]` chunk carries `usage` when448 `stream_options.include_usage` was set, then `data: [DONE]` → done (stop reason from the last finish_reason).449 * Responses API: OpenAI event vocabulary with `sequence_number` — delegated to OpenAINormalizer (incl.450 `response.reasoning_summary_text.delta`, code_interpreter events surfaced as raw).451 * /v1/messages (Anthropic-compatible, DEPRECATED): delegated to AnthropicNormalizer; note live quirks — `index` is 0 for452 both the thinking and the text block and deltas carry no `index` (the delegate tolerates both)."""453454 CHAT_FINISH = {"stop": "end", "length": "max_tokens", "tool_calls": "tool_use", "content_filter": "refusal"}455456 def __init__(self) -> None:457 self.responses = OpenAINormalizer()458 self.messages = AnthropicNormalizer()459 self.chat_id: Optional[str] = None460 self.chat_stop: Optional[str] = None461 self.saw_chat = False462463 def __call__(self, ev: SSEEvent) -> list[Unified]:464 obj = ev.json()465 if not isinstance(obj, dict):466 if ev.data == "[DONE]":467 return [Unified("done", stop_reason=self.chat_stop or "end", raw={"data": "[DONE]"})] if self.saw_chat else []468 return [Unified("raw", raw={"event": ev.event, "data": ev.data})]469 t = obj.get("type") or ev.event or ""470 if obj.get("object") == "chat.completion.chunk" or ("choices" in obj and not t):471 return self._chat(obj)472 if t.startswith("response.") or t in ("error",) and "sequence_number" in obj:473 return self.responses(ev)474 if t in AnthropicNormalizer.MESSAGE_EVENTS:475 return self.messages(ev)476 if t == "error": # bare error object ({"code","error"} shape) — xAI returns HTTP errors before any chunk, but be safe477 return [Unified("error", error=obj, raw=obj)]478 return [Unified("raw", raw=obj)]479480 def _chat(self, obj: dict[str, Any]) -> list[Unified]:481 self.saw_chat = True482 self.chat_id = obj.get("id") or self.chat_id483 out: list[Unified] = []484 choices = obj.get("choices") or []485 if not choices and isinstance(obj.get("usage"), dict):486 return [Unified("usage", usage=obj["usage"], raw=obj)]487 for ch in choices:488 d = ch.get("delta") or {}489 if d.get("reasoning_content"):490 out.append(Unified("reasoning_delta", text=d["reasoning_content"], raw=obj))491 if d.get("content"):492 out.append(Unified("text_delta", text=d["content"], raw=obj))493 for tc in d.get("tool_calls") or []:494 fn = tc.get("function") or {}495 raw_args = fn.get("arguments") or ""496 cid = tc.get("id") or f"{self.chat_id}:{tc.get('index', 0)}"497 out.append(Unified("tool_call_start", tool_call_id=cid, tool_name=fn.get("name"), raw=obj))498 if raw_args:499 out.append(Unified("tool_call_delta", tool_call_id=cid, tool_name=fn.get("name"), arguments_delta=raw_args, raw=obj))500 args = PartialJSON().finish(cid, raw_args or None) if raw_args else {}501 out.append(Unified("tool_call_done", tool_call_id=cid, tool_name=fn.get("name"), arguments=args, raw=obj))502 if ch.get("finish_reason"):503 self.chat_stop = self.CHAT_FINISH.get(ch["finish_reason"], ch["finish_reason"])504 if not out:505 out.append(Unified("raw", raw=obj))506 return out507508509class GeminiNormalizer:510 """Gemini streams → Unified. Two payload families share one normaliser:511 * generateContent chunks (SSE `?alt=sse` or JSON-array framing): `candidates[0].content.parts[]` deltas — `{text}` →512 text_delta, `{text, thought:true}` → reasoning_delta, `{functionCall{name,args,id}}` (+ sibling `thoughtSignature`) →513 tool_call_start/delta/done at once (arguments never arrive partially), empty-text part with `thoughtSignature` → raw514 (keep it: it must be echoed in multi-turn), `inlineData`/`executableCode`/`codeExecutionResult` → raw. `usageMetadata`515 is on EVERY chunk (running totals; the last wins) → one `usage` event when `finishReason` arrives, then `done`516 (STOP→end / tool_use if functionCall seen, MAX_TOKENS→max_tokens, SAFETY & co→refusal). A chunk with `promptFeedback`517 and no candidates = prompt blocked → done(refusal). `{"error": google.rpc.Status}` → error.518 * Interactions API events (`event_type`): interaction.created/status_update → raw; step.start(function_call) →519 tool_call_start; step.delta text → text_delta, thought_summary → reasoning_delta, thought_signature → raw,520 arguments_delta → tool_call_delta; step.stop of a function_call step → tool_call_done; interaction.completed → usage +521 done (status completed→end, requires_action→tool_use, incomplete→incomplete, failed→failed, cancelled→cancelled);522 `error` → error; `event: done` / `data: [DONE]` → nothing (terminal marker after interaction.completed).523 `last_signature` keeps the most recent thoughtSignature for callers that rebuild the model turn by hand."""524525 REFUSAL = {"SAFETY", "RECITATION", "BLOCKLIST", "PROHIBITED_CONTENT", "SPII", "IMAGE_SAFETY", "IMAGE_PROHIBITED_CONTENT",526 "IMAGE_RECITATION", "IMAGE_OTHER", "LANGUAGE"}527 INTERACTION_STATUS = {"completed": "end", "requires_action": "tool_use", "incomplete": "incomplete", "failed": "failed", "cancelled": "cancelled"}528529 def __init__(self) -> None:530 self.args = PartialJSON()531 self.usage: dict[str, Any] = {}532 self.finish_reason: Optional[str] = None533 self.saw_call = False534 self.response_id: Optional[str] = None535 self.model_version: Optional[str] = None536 self.last_signature: Optional[str] = None537 self.steps: dict[int, dict[str, Any]] = {} # Interactions: step index → step header538 self.interaction_id: Optional[str] = None539540 def __call__(self, ev: SSEEvent) -> list[Unified]:541 obj = ev.json()542 if not isinstance(obj, dict):543 return [] if ev.data == "[DONE]" else [Unified("raw", raw={"event": ev.event, "data": ev.data})]544 if "event_type" in obj or ev.event in ("step.delta", "step.start", "step.stop", "interaction.created", "interaction.completed"):545 return self._interaction(obj, obj.get("event_type") or ev.event or "")546 return self._chunk(obj)547548 def _stop(self) -> str:549 fr = self.finish_reason550 if self.saw_call:551 return "tool_use"552 if fr in (None, "STOP"):553 return "end"554 if fr == "MAX_TOKENS":555 return "max_tokens"556 if fr in self.REFUSAL or fr == "BLOCKED":557 return "refusal"558 return "other" # MALFORMED_FUNCTION_CALL, MISSING_THOUGHT_SIGNATURE, UNEXPECTED_TOOL_CALL, TOO_MANY_TOOL_CALLS, MALFORMED_RESPONSE…559560 def _chunk(self, obj: dict[str, Any]) -> list[Unified]:561 u = lambda typ, **kw: Unified(typ, raw=obj, **kw) # noqa: E731562 if isinstance(obj.get("error"), dict):563 return [u("error", error=obj["error"])]564 if isinstance(obj.get("usageMetadata"), dict):565 self.usage = dict(obj["usageMetadata"])566 self.response_id = obj.get("responseId") or self.response_id567 self.model_version = obj.get("modelVersion") or self.model_version568 cands = obj.get("candidates") or []569 out: list[Unified] = []570 if not cands:571 if obj.get("promptFeedback"):572 self.finish_reason = "BLOCKED"573 out.append(u("usage", usage=dict(self.usage)))574 out.append(u("done", stop_reason="refusal"))575 return out576 return [u("raw")]577 c = cands[0]578 for i, p in enumerate((c.get("content") or {}).get("parts") or []):579 if p.get("thoughtSignature"):580 self.last_signature = p["thoughtSignature"]581 if "functionCall" in p:582 fc = p["functionCall"] or {}583 args = fc.get("args") or {}584 cid = fc.get("id") or f"call_{len(self.steps)}_{i}"585 self.saw_call = True586 out.append(u("tool_call_start", tool_call_id=cid, tool_name=fc.get("name")))587 out.append(u("tool_call_delta", tool_call_id=cid, tool_name=fc.get("name"), arguments_delta=json.dumps(args)))588 out.append(u("tool_call_done", tool_call_id=cid, tool_name=fc.get("name"), arguments=args if isinstance(args, dict) else {"_value": args}))589 elif "text" in p:590 if p.get("thought"):591 out.append(u("reasoning_delta", text=p.get("text") or ""))592 elif p.get("text"):593 out.append(u("text_delta", text=p["text"]))594 else:595 out.append(u("raw")) # empty-text signature carrier (Gemini 3 last chunk)596 else:597 out.append(u("raw")) # inlineData, executableCode, codeExecutionResult, fileData…598 if c.get("finishReason"):599 self.finish_reason = c["finishReason"]600 out.append(u("usage", usage=dict(self.usage)))601 out.append(u("done", stop_reason=self._stop()))602 return out or [u("raw")]603604 def _interaction(self, obj: dict[str, Any], t: str) -> list[Unified]:605 u = lambda typ, **kw: [Unified(typ, raw=obj, **kw)] # noqa: E731606 idx = obj.get("index")607 if t == "interaction.created":608 self.interaction_id = (obj.get("interaction") or {}).get("id")609 return u("raw")610 if t == "step.start":611 step = obj.get("step") or {}612 self.steps[idx] = step613 if step.get("type") == "function_call":614 return u("tool_call_start", tool_call_id=step.get("id"), tool_name=step.get("name"))615 return u("raw")616 if t == "step.delta":617 d = obj.get("delta") or {}618 dt = d.get("type")619 step = self.steps.get(idx) or {}620 if dt == "text":621 return u("text_delta", text=d.get("text") or "")622 if dt == "thought_summary":623 return u("reasoning_delta", text=((d.get("content") or {}).get("text") or ""))624 if dt == "thought_signature":625 self.last_signature = d.get("signature") or self.last_signature626 return u("raw")627 if dt == "arguments_delta":628 frag = d.get("arguments") or ""629 self.args.append(str(idx), frag)630 return u("tool_call_delta", tool_call_id=step.get("id"), tool_name=step.get("name"), arguments_delta=frag)631 return u("raw") # image/audio/video/document/function_result/*_call/*_result/processing_*/text_annotation_delta632 if t == "step.stop":633 step = self.steps.get(idx) or {}634 if step.get("type") == "function_call":635 args = self.args.finish(str(idx))636 if not args and isinstance(step.get("arguments"), dict) and step["arguments"]:637 args = step["arguments"]638 return u("tool_call_done", tool_call_id=step.get("id"), tool_name=step.get("name"), arguments=args)639 return u("raw")640 if t == "interaction.completed":641 inter = obj.get("interaction") or {}642 out = []643 if isinstance(inter.get("usage"), dict):644 out += u("usage", usage=inter["usage"])645 return out + u("done", stop_reason=self.INTERACTION_STATUS.get(inter.get("status") or "", inter.get("status") or "end"))646 if t == "error":647 return u("error", error=obj.get("error") or obj)648 return u("raw") # interaction.status_update / in_progress / requires_action and future types649650651def normalizer_for(provider: str):652 return {"openai": OpenAINormalizer, "anthropic": AnthropicNormalizer, "xai": XAINormalizer, "gemini": GeminiNormalizer}[provider]()653654655def iter_unified(chunks: Iterable[bytes], provider: str, framing: str = "sse") -> Iterator[Unified]:656 """chunks: iterable of raw byte chunks from the HTTP body (any chunking). framing: "sse" (all providers) or657 "json_array" (Gemini streamGenerateContent WITHOUT `?alt=sse`)."""658 parser = JSONArrayStreamParser() if framing == "json_array" else SSEParser()659 norm = normalizer_for(provider)660 for chunk in chunks:661 for ev in parser.feed(chunk):662 yield from norm(ev)663 for ev in parser.close():664 yield from norm(ev)665666667def collect_text(events: Iterable[Unified]) -> str:668 return "".join(e.text for e in events if e.type == "text_delta")669670671if __name__ == "__main__":672 import sys673 if len(sys.argv) not in (3, 4):674 print("usage: sse_parser.py <openai|anthropic|xai|gemini> <recorded.sse|recorded.json> [sse|json_array]")675 sys.exit(1)676 data = open(sys.argv[2], "rb").read()677 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"):678 if u.type != "raw":679 print(u.type, u.text or u.arguments_delta or u.arguments or u.usage or u.stop_reason or u.error or "")680