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.9 KB · 680 lines python
Raw Blame History
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