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%
19.0 KB · 318 lines python
Raw Blame History
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