Python 88.3%
TypeScript 7.6%
Shell 4.1%
1"""Offline tests for examples/shared/streaming/sse_parser.py using recorded fixtures in tests/shared/fixtures/sse/."""2from __future__ import annotations34import json5import sys6from pathlib import Path78import pytest910ROOT = Path(__file__).resolve().parents[2]11sys.path.insert(0, str(ROOT / "examples" / "shared" / "streaming"))12FIX = ROOT / "tests" / "shared" / "fixtures" / "sse"1314import sse_parser as sp # noqa: E4021516OPENAI = (FIX / "openai-responses-tool-call.sse").read_bytes()17ANTHROPIC = (FIX / "anthropic-messages-tool-use.sse").read_bytes()18OVERLOADED = (FIX / "anthropic-messages-overloaded.sse").read_bytes()192021# ---------------- generic framing ----------------2223@pytest.mark.parametrize("chunk", [None, 1, 7, 64, 4096])24def test_chunk_boundaries_do_not_change_events(chunk):25 evs = sp.parse_sse_bytes(OPENAI, chunk)26 assert len(evs) == 16 and evs[0].event == "response.created" and evs[-1].event == "response.completed"27 assert all(e.json() is not None for e in evs)282930def test_comments_multiline_data_crlf_bom_and_done_sentinel():31 raw = b"\xef\xbb\xbf: hello\r\nevent: x\r\ndata: {\"a\":\r\ndata: 1}\r\nid: 7\r\nretry: 3000\r\n\r\ndata: [DONE]\r\n\r\n"32 evs = sp.parse_sse_bytes(raw)33 assert len(evs) == 234 assert evs[0].event == "x" and evs[0].data == '{"a":\n1}' and evs[0].json() == {"a": 1} and evs[0].id == "7" and evs[0].retry == 300035 assert evs[1].event is None and evs[1].data == "[DONE]"363738def test_trailing_event_without_blank_line_is_flushed_on_close():39 evs = sp.parse_sse_bytes(b"event: message_stop\ndata: {\"type\":\"message_stop\"}")40 assert len(evs) == 1 and evs[0].json()["type"] == "message_stop"414243def test_data_without_leading_space_and_unknown_fields():44 evs = sp.parse_sse_bytes(b"data:{\"k\":1}\nfoo: bar\n\n")45 assert evs[0].json() == {"k": 1}464748def test_buffer_overflow_guard():49 p = sp.SSEParser(max_buffer=100)50 with pytest.raises(sp.SSEBufferOverflow):51 p.feed(b"data: " + b"x" * 200)525354def test_incremental_feed_returns_events_as_soon_as_complete():55 p = sp.SSEParser()56 assert p.feed(b"event: a\ndata: {\"t\":1}\n") == []57 got = p.feed(b"\nevent: b\n")58 assert len(got) == 1 and got[0].event == "a"59 assert p.feed(b"data: {}\n\n")[0].event == "b"606162# ---------------- partial JSON ----------------6364@pytest.mark.parametrize("prefix,expected", [65 ('{"loc', {}),66 ('{"location":"Par', {"location": "Par"}),67 ('{"location":"Paris","unit":', {"location": "Paris"}),68 ('{"location":"Paris","unit":"c"}', {"location": "Paris", "unit": "c"}),69 ('{"a":[1,2', {"a": [1, 2]}),70 ('{"a":{"b":"x', {"a": {"b": "x"}}),71 ("", None),72])73def test_parse_partial_json(prefix, expected):74 assert sp.parse_partial_json(prefix) == expected757677def test_partial_json_accumulator_preview_and_finish():78 acc = sp.PartialJSON()79 acc.append("fc_1", '{"city":"Pa')80 assert acc.preview("fc_1") == {"city": "Pa"}81 acc.append("fc_1", 'ris"}')82 assert acc.finish("fc_1") == {"city": "Paris"} and "fc_1" not in acc.buffers83 acc.append("bad", "{oops")84 assert acc.finish("bad")["_error"] == "invalid JSON"85 assert acc.finish("missing") == {}868788# ---------------- OpenAI normalisation ----------------8990def test_openai_fixture_normalised():91 evs = list(sp.iter_unified([OPENAI[i:i + 33] for i in range(0, len(OPENAI), 33)], "openai"))92 assert sp.collect_text(evs) == "Let me check the weather."93 types = [e.type for e in evs if e.type != "raw"]94 assert types == ["text_delta", "text_delta", "tool_call_start", "tool_call_delta", "tool_call_delta", "tool_call_delta", "tool_call_done", "usage", "done"]95 start = next(e for e in evs if e.type == "tool_call_start")96 assert start.tool_call_id == "call_abc123" and start.tool_name == "get_weather"97 deltas = "".join(e.arguments_delta for e in evs if e.type == "tool_call_delta")98 assert deltas == '{"location":"Paris","unit":"c"}'99 done = next(e for e in evs if e.type == "tool_call_done")100 assert done.tool_call_id == "call_abc123" and done.arguments == {"location": "Paris", "unit": "c"}101 usage = next(e for e in evs if e.type == "usage").usage102 assert usage["input_tokens"] == 36 and usage["total_tokens"] == 64103 assert evs[-1].type == "done" and evs[-1].stop_reason == "end" and evs[-1].sequence == 15104105106def test_openai_sequence_cursor_tracked_for_resume():107 parser, norm = sp.SSEParser(), sp.OpenAINormalizer()108 for ev in parser.feed(OPENAI[:1500]):109 norm(ev)110 assert norm.response_id == "resp_fixture_1" and norm.last_sequence is not None and 0 < norm.last_sequence < 15111112113def test_openai_incomplete_and_error_events():114 raw = (b'event: response.incomplete\ndata: {"type":"response.incomplete","sequence_number":3,"response":{"status":"incomplete","incomplete_details":{"reason":"max_output_tokens"},"usage":{"output_tokens":16}}}\n\n'115 b'event: error\ndata: {"type":"error","code":"server_error","message":"boom","sequence_number":4}\n\n')116 evs = list(sp.iter_unified([raw], "openai"))117 assert [e.type for e in evs] == ["usage", "done", "error"] and evs[1].stop_reason == "max_tokens" and evs[2].error["code"] == "server_error"118119120# ---------------- Anthropic normalisation ----------------121122def test_anthropic_fixture_normalised():123 evs = list(sp.iter_unified([ANTHROPIC[i:i + 50] for i in range(0, len(ANTHROPIC), 50)], "anthropic"))124 assert sp.collect_text(evs) == "Okay, let me check the weather."125 types = [e.type for e in evs if e.type != "raw"]126 assert types == ["text_delta", "text_delta", "tool_call_start"] + ["tool_call_delta"] * 7 + ["tool_call_done", "usage", "done"]127 done = next(e for e in evs if e.type == "tool_call_done")128 assert done.tool_call_id == "toolu_01T1x1fJ34qAmk2tNTrN7Up6" and done.tool_name == "get_weather"129 assert done.arguments == {"location": "San Francisco, CA", "unit": "fahrenheit"}130 usage = next(e for e in evs if e.type == "usage").usage131 assert usage["input_tokens"] == 472 and usage["output_tokens"] == 89 # message_start input + message_delta output merged132 assert evs[-1].type == "done" and evs[-1].stop_reason == "tool_use"133 assert sum(1 for e in evs if e.type == "raw" and e.raw.get("type") == "ping") == 1 # ping surfaced as raw, not dropped134135136def test_anthropic_in_stream_overloaded_error():137 evs = list(sp.iter_unified([OVERLOADED], "anthropic"))138 assert sp.collect_text(evs) == "O"139 assert evs[-1].type == "error" and evs[-1].error["type"] == "overloaded_error"140 assert not any(e.type == "done" for e in evs) # consumer must treat as failed/restart141142143def test_anthropic_partial_preview_mid_stream():144 parser, norm = sp.SSEParser(), sp.AnthropicNormalizer()145 cut = ANTHROPIC.find(b'" CA\\", "') # stop before the unit key arrives146 for ev in parser.feed(ANTHROPIC[:cut]):147 norm(ev)148 assert norm.args.preview("1") == {"location": "San Francisco,"}149150151def test_unknown_event_types_are_raw_not_errors():152 evs = list(sp.iter_unified([b'event: future_event\ndata: {"type":"future_event","x":1}\n\n'], "anthropic"))153 assert evs[0].type == "raw" and evs[0].raw == {"type": "future_event", "x": 1}154 evs = list(sp.iter_unified([b'event: response.something_new\ndata: {"type":"response.something_new","sequence_number":9}\n\n'], "openai"))155 assert evs[0].type == "raw" and evs[0].sequence == 9156157158# ---------------- xAI normalisation (fixtures cut from tmp-live/xai/*_stream*.json, 2026-09-19) ----------------159160XAI_CHAT = (FIX / "xai-chat-reasoning-tool-call.sse").read_bytes()161XAI_RESP = (FIX / "xai-responses-reasoning-function.sse").read_bytes()162163164@pytest.mark.parametrize("chunk", [1, 17, 4096])165def test_xai_chat_chunks_reasoning_whole_tool_call_usage_done(chunk):166 evs = list(sp.iter_unified([XAI_CHAT[i:i + chunk] for i in range(0, len(XAI_CHAT), chunk)], "xai"))167 types = [e.type for e in evs if e.type != "raw"]168 assert types == ["reasoning_delta", "reasoning_delta", "reasoning_delta", "tool_call_start", "tool_call_delta", "tool_call_done", "usage", "done"]169 assert "".join(e.text for e in evs if e.type == "reasoning_delta") == "I'll use the get_weather tool to check Paris's weather."170 done = next(e for e in evs if e.type == "tool_call_done")171 assert done.tool_call_id == "call-755df5ce-34fb-4354-b73d-21197362d8af-0" and done.tool_name == "get_weather" and done.arguments == {"city": "Paris"}172 usage = next(e for e in evs if e.type == "usage").usage173 assert usage["prompt_tokens"] == 196 and usage["completion_tokens_details"]["reasoning_tokens"] == 121 and usage["prompt_tokens_details"]["cached_tokens"] == 192174 assert evs[-1].type == "done" and evs[-1].stop_reason == "tool_use" and evs[-1].raw == {"data": "[DONE]"}175 assert not any(e.type == "raw" and e.raw.get("data") == "[DONE]" for e in evs) # sentinel consumed, not surfaced as raw176177178def test_xai_chat_plain_text_stream_ends_with_end():179 raw = (b'data: {"id":"d8","object":"chat.completion.chunk","created":0,"model":"grok-4.3","choices":[{"index":0,"delta":{"content":"391","role":"assistant"}}]}\n\n'180 b'data: {"id":"d8","object":"chat.completion.chunk","created":0,"model":"grok-4.3","choices":[{"index":0,"delta":{},"finish_reason":"stop"}]}\n\n'181 b'data: [DONE]\n\n')182 evs = list(sp.iter_unified([raw], "xai"))183 assert [e.type for e in evs] == ["text_delta", "raw", "done"] and sp.collect_text(evs) == "391" and evs[-1].stop_reason == "end"184185186def test_xai_responses_events_delegate_to_openai_vocabulary():187 evs = list(sp.iter_unified([XAI_RESP[i:i + 41] for i in range(0, len(XAI_RESP), 41)], "xai"))188 types = [e.type for e in evs if e.type != "raw"]189 assert types == ["reasoning_delta", "reasoning_delta", "reasoning_delta", "tool_call_start", "tool_call_delta", "tool_call_done", "usage", "done"]190 assert "".join(e.text for e in evs if e.type == "reasoning_delta") == "I need to calculate 17 times 23.\n\n17 times 23 equals 391. I will call get_weather for Paris."191 start = next(e for e in evs if e.type == "tool_call_start")192 assert start.tool_call_id == "call-755df5ce-34fb-4354-b73d-21197362d8af-0" and start.tool_name == "get_weather"193 delta = next(e for e in evs if e.type == "tool_call_delta")194 assert delta.arguments_delta == '{"city":"Paris"}' # xAI sends the whole JSON in ONE delta195 done = next(e for e in evs if e.type == "tool_call_done")196 assert done.arguments == {"city": "Paris"} and done.tool_call_id == "call-755df5ce-34fb-4354-b73d-21197362d8af-0"197 usage = next(e for e in evs if e.type == "usage").usage198 assert usage["output_tokens"] == 151 and usage["output_tokens_details"]["reasoning_tokens"] == 133199 assert evs[-1].type == "done" and evs[-1].stop_reason == "end" and evs[-1].sequence == 16200 norm = sp.XAINormalizer()201 for ev in sp.parse_sse_bytes(XAI_RESP[:2000]):202 norm(ev)203 assert norm.responses.response_id == "ee2a3737-1808-925b-8c71-f5b04b61b711" and norm.responses.last_sequence is not None # cursor kept (no resume on xAI, but useful for logs)204205206def test_xai_messages_compat_events_delegate_to_anthropic_vocabulary():207 raw = (b'event: message_start\ndata: {"type":"message_start","message":{"id":"1bf6522c","type":"message","role":"assistant","content":[],"model":"grok-4.3","usage":{"input_tokens":4,"cache_read_input_tokens":192,"output_tokens":0}}}\n\n'208 b'event: content_block_start\ndata: {"type":"content_block_start","index":0,"content_block":{"type":"thinking","signature":"","thinking":""}}\n\n'209 b'event: content_block_delta\ndata: {"type":"content_block_delta","delta":{"type":"thinking_delta","thinking":"The user requested a"}}\n\n'210 b'event: content_block_stop\ndata: {"type":"content_block_stop","index":0}\n\n'211 b'event: content_block_start\ndata: {"type":"content_block_start","index":0,"content_block":{"type":"text","text":""}}\n\n'212 b'event: content_block_delta\ndata: {"type":"content_block_delta","delta":{"type":"text_delta","text":"OK."}}\n\n'213 b'event: content_block_stop\ndata: {"type":"content_block_stop","index":0}\n\n'214 b'event: message_delta\ndata: {"type":"message_delta","delta":{"stop_reason":"end_turn","stop_sequence":null},"usage":{"output_tokens":147}}\n\n'215 b'event: message_stop\ndata: {"type":"message_stop"}\n\n')216 evs = list(sp.iter_unified([raw], "xai"))217 types = [e.type for e in evs if e.type != "raw"]218 assert types == ["reasoning_delta", "text_delta", "usage", "done"] # index reused (0 twice) and deltas without index tolerated219 assert sp.collect_text(evs) == "OK." and evs[-1].stop_reason == "end"220 assert next(e for e in evs if e.type == "usage").usage["output_tokens"] == 147221222223# ---------------- Gemini normalisation (fixtures cut from tmp-live/gemini-core/*.raw.txt and gemini-tools/h2_interaction_stream.json) ----------------224225GEM_SSE = (FIX / "gemini-stream-thought-signature.sse").read_bytes()226GEM_ARRAY = (FIX / "gemini-stream-json-array.json").read_bytes()227GEM_INTER = (FIX / "gemini-interactions-step-delta.sse").read_bytes()228GEM_INTER_FC = (FIX / "gemini-interactions-function-call.sse").read_bytes()229230231@pytest.mark.parametrize("chunk", [1, 13, 4096])232def test_gemini_alt_sse_thought_text_signature_usage(chunk):233 evs = list(sp.iter_unified([GEM_SSE[i:i + chunk] for i in range(0, len(GEM_SSE), chunk)], "gemini"))234 types = [e.type for e in evs]235 assert types == ["reasoning_delta", "text_delta", "text_delta", "text_delta", "raw", "usage", "done"]236 assert sp.collect_text(evs) == "1 2 3 4 5 6 7 8 9 10 11 12" and evs[0].text.startswith("**Counting**")237 assert evs[4].raw["candidates"][0]["content"]["parts"][0]["thoughtSignature"].startswith("El4K") # empty-text signature part kept as raw238 usage = evs[5].usage239 assert usage["promptTokenCount"] == 13 and usage["candidatesTokenCount"] == 26 and usage["thoughtsTokenCount"] == 40 # LAST chunk's usage240 assert evs[-1].stop_reason == "end"241242243def test_gemini_normalizer_state_response_id_and_signature():244 norm = sp.GeminiNormalizer()245 for ev in sp.parse_sse_bytes(GEM_SSE):246 norm(ev)247 assert norm.response_id == "bwWuasWFKLCg_PUP1deUqA4" and norm.model_version == "gemini-3.5-flash-lite"248 assert norm.last_signature.startswith("El4KXAFpFH0Tj8ibw628") and norm.finish_reason == "STOP"249250251@pytest.mark.parametrize("chunk", [1, 7, 100, 100000])252def test_gemini_json_array_framing_same_events_as_sse(chunk):253 evs = list(sp.iter_unified([GEM_ARRAY[i:i + chunk] for i in range(0, len(GEM_ARRAY), chunk)], "gemini", framing="json_array"))254 assert [e.type for e in evs] == ["text_delta", "text_delta", "text_delta", "raw", "usage", "done"]255 assert sp.collect_text(evs) == "1 2 3 4 5 6 7 8 9 10 11 12" and evs[-1].stop_reason == "end"256 assert evs[4].raw["candidates"][0]["content"]["parts"][0]["thoughtSignature"].endswith('"}{"}') # braces inside strings do not confuse the parser257258259def test_json_array_parser_emits_objects_incrementally_and_bounds_memory():260 p = sp.JSONArrayStreamParser()261 assert p.feed(b"[{\n \"a\": 1") == []262 got = p.feed(b"\n}\n,\n{\"b\": \"}\"")263 assert len(got) == 1 and got[0].json() == {"a": 1}264 assert p.feed(b"}\n]")[0].json() == {"b": "}"} and p.close() == []265 with pytest.raises(sp.SSEBufferOverflow):266 sp.JSONArrayStreamParser(max_buffer=50).feed(b"[{" + b'"k":"' + b"x" * 100)267268269def test_gemini_function_call_chunk_and_prompt_block_and_error():270 fc = json.dumps({"candidates": [{"content": {"parts": [{"functionCall": {"name": "get_weather", "args": {"city": "Paris"}, "id": "call_1"}, "thoughtSignature": "SIG"}], "role": "model"}, "finishReason": "STOP", "index": 0}],271 "usageMetadata": {"promptTokenCount": 30, "candidatesTokenCount": 10, "totalTokenCount": 40}}).encode()272 evs = list(sp.iter_unified([b"data: " + fc + b"\r\n\r\n"], "gemini"))273 assert [e.type for e in evs] == ["tool_call_start", "tool_call_delta", "tool_call_done", "usage", "done"]274 assert evs[2].tool_call_id == "call_1" and evs[2].arguments == {"city": "Paris"} and evs[1].arguments_delta == '{"city": "Paris"}' and evs[-1].stop_reason == "tool_use"275 blocked = json.dumps({"promptFeedback": {"blockReason": "PROHIBITED_CONTENT"}, "usageMetadata": {"promptTokenCount": 9, "totalTokenCount": 9}}).encode()276 evs = list(sp.iter_unified([b"data: " + blocked + b"\n\n"], "gemini"))277 assert [e.type for e in evs] == ["usage", "done"] and evs[-1].stop_reason == "refusal"278 err = json.dumps({"error": {"code": 429, "message": "You exceeded your current quota… Please retry in 54.22s.", "status": "RESOURCE_EXHAUSTED"}}).encode()279 evs = list(sp.iter_unified([b"data: " + err + b"\n\n"], "gemini"))280 assert evs[0].type == "error" and evs[0].error["status"] == "RESOURCE_EXHAUSTED"281 trunc = json.dumps({"candidates": [{"content": {"parts": [{"text": "par"}], "role": "model"}, "index": 0}]}).encode()282 evs = list(sp.iter_unified([b"data: " + trunc + b"\n\n"], "gemini"))283 assert [e.type for e in evs] == ["text_delta"] # no finishReason → no done: the consumer must treat the close as an error284285286def test_gemini_max_tokens_finish_reason():287 raw = b'data: {"candidates": [{"content": {"parts": [{"text": "{\\"ok\\": true, \\"word\\": \\"OK"}],"role": "model"},"finishReason": "MAX_TOKENS","index": 0}],"usageMetadata": {"promptTokenCount": 5,"candidatesTokenCount": 16,"totalTokenCount": 21}}\r\n\r\n'288 evs = list(sp.iter_unified([raw], "gemini"))289 assert [e.type for e in evs] == ["text_delta", "usage", "done"] and evs[-1].stop_reason == "max_tokens"290291292def test_gemini_interactions_step_events_text_and_signature():293 evs = list(sp.iter_unified([GEM_INTER[i:i + 29] for i in range(0, len(GEM_INTER), 29)], "gemini"))294 types = [e.type for e in evs]295 assert types == ["raw", "raw", "raw", "raw", "raw", "raw", "text_delta", "raw", "usage", "done"]296 assert sp.collect_text(evs) == "OK." and evs[-1].stop_reason == "end"297 assert evs[-2].usage["total_input_tokens"] == 13 and evs[-2].usage["total_output_tokens"] == 2298 norm = sp.GeminiNormalizer()299 for ev in sp.parse_sse_bytes(GEM_INTER):300 norm(ev)301 assert norm.interaction_id.startswith("v1_Chd") and norm.last_signature.startswith("El4KXAFpFH0TGIU")302 assert not any(e.type == "raw" and e.raw.get("data") == "[DONE]" for e in evs) # `event: done` / [DONE] marker swallowed303304305def test_gemini_interactions_function_call_steps_requires_action():306 evs = list(sp.iter_unified([GEM_INTER_FC], "gemini"))307 types = [e.type for e in evs if e.type != "raw"]308 assert types == ["reasoning_delta", "tool_call_start", "tool_call_delta", "tool_call_delta", "tool_call_done", "usage", "done"]309 assert evs[[e.type for e in evs].index("reasoning_delta")].text == "I need to find the weather in San Francisco."310 done = next(e for e in evs if e.type == "tool_call_done")311 assert done.tool_call_id == "call_8f2c" and done.tool_name == "get_weather" and done.arguments == {"location": "San Francisco, CA"}312 assert evs[-1].stop_reason == "tool_use" # status requires_action → the client must POST function_result + previous_interaction_id313314315def test_normalizer_for_four_providers():316 for name, cls in (("openai", sp.OpenAINormalizer), ("anthropic", sp.AnthropicNormalizer), ("xai", sp.XAINormalizer), ("gemini", sp.GeminiNormalizer)):317 assert isinstance(sp.normalizer_for(name), cls)318