# Streaming patterns — one SSE parser for OpenAI Responses, Anthropic Messages, xAI (3 framings) and Gemini (SSE, JSON array, Interactions) **Status:** DOCUMENTED (framing + event vocabularies) · offline-tested on recorded fixtures (`tests/shared/test_sse_parser.py`, 45 tests; `sseParser.ts --selftest`) · text streaming LIVE_VERIFIED through the provider adapters (OpenAI/Anthropic 2026-09-18; xAI Responses + Chat and Gemini `?alt=sse` 2026-09-19) **Sources:** - OpenAI OpenAPI spec (`response.*` events, `sequence_number`), https://developers.openai.com/api/docs/guides/background - https://platform.claude.com/docs/en/build-with-claude/streaming - xAI: https://docs.x.ai/developers/model-capabilities/text/streaming · …/text/reasoning (summary deltas) · https://docs.x.ai/developers/tools/streaming-and-sync · `generated/fragments/streaming-events/xai-inference.json` (5 live captures in `tmp-live/xai/*_stream*.json`) - Gemini: https://ai.google.dev/api/generate-content#method:-models.streamgeneratecontent · https://ai.google.dev/gemini-api/docs/thought-signatures · https://ai.google.dev/api/interactions (SSE events) · https://ai.google.dev/api/live · `generated/fragments/streaming-events/{gemini-core,gemini-interactions,gemini-live,gemini-lyria-realtime}.json` (captures in `tmp-live/gemini-core/*.raw.txt`, `tmp-live/gemini-tools/h2_interaction_stream.json`) - WHATWG HTML §9.2 Server-sent events (framing) **Last verified:** 2026-09-19 Implementation: `examples/shared/streaming/sse_parser.py`, `sseParser.ts`. Fixtures: `tests/shared/fixtures/sse/*` — `openai-responses-tool-call.sse`, `anthropic-messages-tool-use.sse`, `anthropic-messages-overloaded.sse`, `xai-chat-reasoning-tool-call.sse`, `xai-responses-reasoning-function.sse`, `gemini-stream-thought-signature.sse`, `gemini-stream-json-array.json`, `gemini-interactions-step-delta.sse`, `gemini-interactions-function-call.sse` (xAI/Gemini fixtures cut from the live captures; `thoughtSignature`/`encrypted_content` blobs truncated). ## 1. Wire framing — six framings, one parser | | OpenAI Responses | Anthropic Messages | xAI Chat Completions | xAI Responses | Gemini `streamGenerateContent?alt=sse` | Gemini `streamGenerateContent` (default) | Gemini Interactions (`stream: true`) | |---|---|---|---|---|---|---|---| | Content-Type | `text/event-stream` | `text/event-stream` | `text/event-stream` | `text/event-stream` | `text/event-stream` | **`application/json`** | `text/event-stream` | | Lines | `event:` + `data:` + blank | `event:` + `data:` + blank | **`data:` only** (no `event:`) | `event:` + `data:` + blank | **`data:` only**, `\r\n\r\n` | pretty-printed **JSON array** streamed incrementally: `[{`, objects separated by a line `,`, closed by `]` | `event:` + `data:` | | Authoritative type | JSON `type` | JSON `type` | `object: "chat.completion.chunk"` + `choices[].delta` | JSON `type` | payload shape (`candidates[]`) | same | JSON `event_type` | | Ordering / cursor | `sequence_number` | block `index` | chunk order | `sequence_number` (no resume: `background` rejected) | chunk order; `responseId` constant | same | `index` per step; `last_event_id` resume for `background` runs | | Keep-alive | none | `event: ping` | none | none | none | none | none documented | | Terminal | `response.completed/incomplete/failed/cancelled` | `message_stop` | `data: [DONE]` (usage chunk with `choices: []` before it when `stream_options.include_usage`) | `response.completed` (no `[DONE]`) | last chunk carries `finishReason` (+ empty-text `thoughtSignature` part on Gemini 3) | same | `interaction.completed` then `event: done` / `data: [DONE]` | | Errors after 200 | `event: error` | `event: error` (`overloaded_error`) | none observed (errors precede the first chunk) | `response.failed` / `error` (not observed) | not documented; `{"error": google.rpc}` chunk handled defensively; a close without `finishReason` = failure | same | `event: error` (`{error:{code:"", message}}`) | The generic `SSEParser` implements the spec strictly (`event:`, multi-line `data:`, `id:`, `retry:`, comments, CRLF/LF, BOM, unknown fields ignored, dispatch on blank line, flush on close), is **incremental** (any chunk boundaries — fixtures run at 1-, 7-, 13-, 17-, 29-, 33-, 41-, 50-, 64- and 4096-byte chunks) and **bounded** (`max_buffer`, 8 MiB). `JSONArrayStreamParser` does the same for Gemini's non-SSE mode: it tracks string/escape state and brace depth, emits each top-level object as soon as it closes (braces inside strings — e.g. inside a `thoughtSignature` — do not confuse it), and never buffers more than one object. Prefer `?alt=sse` (the SDKs always use it); the array parser exists for proxies and logs that recorded the default framing. ## 2. Backpressure - Python: `iter_unified(chunks, provider, framing="sse"|"json_array")` is a generator; it pulls the next chunk only when the consumer asks for the next event. - TypeScript: `iterUnified(body, provider, framing)` reads the `ReadableStream` with a reader, one `read()` per consumed chunk. Never `await r.text()` a stream. - Cancellation: abort the fetch / close the socket; OpenAI background responses `POST /v1/responses/{id}/cancel`; Gemini Interactions `POST /v1beta/interactions/{id}/cancel`; xAI has no cancel (close the connection; billed for what was generated). ## 3. Reasoning in streams | Provider | Event | Unified | |---|---|---| | OpenAI | `response.reasoning_summary_text.delta`, `response.reasoning_text.delta` | `reasoning_delta` | | Anthropic | `content_block_delta` / `thinking_delta` (+ `signature_delta` → raw) | `reasoning_delta` | | xAI Chat | `choices[].delta.reasoning_content` — arrives **before** `content`; absent when the model produces no summary ("Reply with OK." had none) | `reasoning_delta` | | xAI Responses | `response.reasoning_summary_text.delta` (emitted on grok-4.3 without any `reasoning.summary` setting); `response.reasoning_text.delta` documented, not observed; the reasoning `output_item.done` may carry `encrypted_content` | `reasoning_delta` | | xAI `/v1/messages` (compat) | `thinking_delta` in a block with `index: 0` — the following text block **also** has `index: 0` and deltas carry no `index` (live quirk; the normaliser tolerates it) | `reasoning_delta` | | Gemini | parts `{text, thought: true}` (rolling summaries, only with `thinkingConfig.includeThoughts`; not guaranteed even when `thoughtsTokenCount` is billed) | `reasoning_delta` | | Gemini Interactions | `step.start {step:{type:"thought"}}` → `step.delta {type:"thought_summary", content:{text}}` / `{type:"thought_signature", signature}` | `reasoning_delta` / raw (signature kept in `last_signature`) | **Thought signatures (Gemini 3):** the last chunk of every stream carries a part with **empty text and `thoughtSignature`**; `functionCall` parts carry their own. The normaliser surfaces them as `raw` (never dropped) and remembers `last_signature`; in multi-turn you must echo the model turn verbatim (`docs/gemini/tool-loop.md`). Parsers that skip empty text parts break Gemini 3 tool calling. ## 4. Partial JSON for tool arguments | | OpenAI | Anthropic | xAI Chat | xAI Responses | Gemini | Gemini Interactions | |---|---|---|---|---|---|---| | Start | `output_item.added` (`function_call`) | `content_block_start` (`tool_use`) | the `tool_calls[]` chunk itself | `output_item.added` (`function_call`, `call_id: "call-…"`) | the `functionCall` part | `step.start {type:"function_call", id, name}` | | Fragments | `function_call_arguments.delta` | `input_json_delta.partial_json` | **none — the whole call (id, name, arguments) arrives in ONE chunk** | one `function_call_arguments.delta` containing the **whole** JSON string (live) | **none — `args` is a complete object** | `step.delta {type:"arguments_delta", arguments}` (partial JSON string; accumulate) | | End | `function_call_arguments.done` (authoritative `arguments`) | `content_block_stop` | `finish_reason: "tool_calls"` | `function_call_arguments.done` (`arguments`, `name`) | same part | `step.stop`; `interaction.completed {status: "requires_action"}` | | Keyed by | `item_id` → `call_id` | block `index` → `tool_use.id` | `tool_calls[].id` (`index` fallback) | `item_id` → `call_id` | `functionCall.id` (or synthetic) | step `index` → `step.id` (`call_…`) | `PartialJSON` accumulates per key and `parse_partial_json()` offers a *preview*; the authoritative value is always the `done`/stop payload. For xAI and Gemini generateContent the normaliser still emits `tool_call_start → tool_call_delta → tool_call_done` so consumers have one code path — the "delta" is simply complete. ## 5. Unified event vocabulary | Unified | OpenAI | Anthropic | xAI Chat | xAI Responses | Gemini generateContent | Gemini Interactions | |---|---|---|---|---|---|---| | `text_delta` | `output_text.delta`, `refusal.delta` | `text_delta` | `delta.content` | `output_text.delta` | non-thought `text` parts (non-empty) | `step.delta {type:"text"}` | | `reasoning_delta` | reasoning summary/text deltas | `thinking_delta` | `delta.reasoning_content` | reasoning summary deltas | `thought: true` parts | `thought_summary` | | `tool_call_*` | function/custom/mcp call items | `tool_use`/`server_tool_use`/`mcp_tool_use` blocks | `delta.tool_calls[]` | same as OpenAI (x_search appears as `custom_tool_call`) | `functionCall` parts | `function_call` steps | | `usage` | terminal event `response.usage` | `message_start` ⊕ `message_delta` | final `choices: []` chunk (`stream_options.include_usage`) — `completion_tokens` excludes reasoning | terminal `response.usage` (`output_tokens` includes reasoning, `cost_in_usd_ticks`) | `usageMetadata` of the **last** chunk (running totals on every chunk; `thoughtsTokenCount`) | `interaction.completed.interaction.usage` (`total_input_tokens`, `total_output_tokens`, `total_thought_tokens`…) | | `done` | terminal event | `message_stop` | `[DONE]`; stop from last `finish_reason` (`stop`/`length`/`tool_calls`/`content_filter`) | `response.completed` | on `finishReason` (`STOP`→end / `tool_use`, `MAX_TOKENS`, safety enums→`refusal`, tool-call errors→`other`); `promptFeedback` without candidates → `refusal` | `interaction.completed.status` (`completed`→end, `requires_action`→tool_use, `incomplete`, `failed`, `cancelled`) | | `error` | `event: error` | `event: error` | — | `error` | `{"error": …}` chunk | `event: error` | | `raw` | `response.created`, `content_part.*`, code_interpreter events… | `message_start`, `ping`, `signature_delta` | role-only / empty deltas | `response.in_progress`, `output_item.done`, `reasoning_summary_part.*` | empty-text signature parts, `inlineData`, `executableCode`, `codeExecutionResult` | `interaction.created`, `status_update`, `step.start/stop` of non-tool steps, `thought_signature`, media/grounding deltas | `XAINormalizer` dispatches on payload shape (`chat.completion.chunk` / `response.*` / Anthropic event names) so one instance handles any xAI endpoint; `GeminiNormalizer` handles both generateContent chunks and Interactions events (`event_type`). Unknown types are surfaced as `raw`, never raised. ## 6. WebSocket surfaces (not SSE — normalisation notes only) - **Gemini Live API** (`wss://…/BidiGenerateContent`, `docs/gemini/live-events.md`): one JSON message per frame; server messages are a flattened union (`setupComplete` | `serverContent{modelTurn, generationComplete, interrupted, turnComplete, inputTranscription, outputTranscription}` | `toolCall{functionCalls[]}` | `toolCallCancellation{ids}` | `goAway` | `sessionResumptionUpdate`) plus periodic `usageMetadata`. Mapping: `serverContent.modelTurn.parts[]` → `text_delta` (or audio `inlineData` → raw), `toolCall.functionCalls[]` → `tool_call_start/done` (arguments whole; reply with `toolResponse`), `turnComplete` → `done`, `goAway` → reconnect with `sessionResumption.handle` (valid 2 h; connection resets ≈ 10 min; audio session 15 min / audio+video 2 min without `contextWindowCompression`), `toolCallCancellation` → discard in-flight tool work (user barge-in). No SSE framing: feed each frame's JSON to a normaliser directly. - **Gemini Lyria RealTime** (`BidiGenerateMusic`): `setup` → `setupComplete`; `serverContent.audioChunks[]` (raw 16-bit PCM 48 kHz stereo, faster than real time → client buffers), `filteredPrompt` (safety-rejected prompt). Only "audio delta" and "filtered" events exist; no text/tool vocabulary. - **xAI**: `wss://api.x.ai/v1/responses` (Responses over WebSocket, same event objects as the SSE stream), `wss://api.x.ai/v1/realtime` (speech-to-speech, OpenAI-Realtime-like events), `/v1/stt`, `/v1/tts` — documented, not exercised by this parser. ## 7. Reconnection (see `resilience.md` §7) - OpenAI: `last_sequence` + `response_id` → `GET /v1/responses/{id}?stream=true&starting_after=`. xAI: same cursor but **no resume** (`background` → 400) — restart. - Anthropic: restart with a `restart` marker. Gemini generateContent: restart (no cursor); Gemini Interactions: `GET …/interactions/{id}?stream=true&last_event_id=…` for `background: true`; Live: `sessionResumption`. ## 8. Checklist - [ ] Parse by blank-line dispatch (SSE) or brace depth (Gemini JSON array); never "one JSON per line". - [ ] Dispatch on the JSON payload (`type` / `event_type` / shape), tolerate unknown types, swallow `[DONE]` sentinels (xAI chat, Gemini Interactions). - [ ] Accumulate tool arguments as strings where partial (OpenAI, Anthropic, Interactions); treat xAI and Gemini whole-argument deliveries through the same start/delta/done path. - [ ] Never drop empty text parts on Gemini — they carry `thoughtSignature`; keep `last_signature` and echo model turns verbatim. - [ ] Usage: merge Anthropic `message_start` + `message_delta`; take Gemini's **last** `usageMetadata`; request `stream_options.include_usage` on xAI chat; read `output_tokens_details.reasoning_tokens` / `thoughtsTokenCount` to explain the bill. - [ ] Treat `event: error` after HTTP 200 (Anthropic, Interactions) and a Gemini close without `finishReason` as failures. - [ ] Bound buffers; stream with backpressure; never `.text()` a stream. - [ ] Keep `sequence_number` for OpenAI resume; for xAI it is diagnostics only.