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%
14.0 KB

# 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:

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:"<snake_case>", 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.