From 3ce1d96f4ddac1060d9b5e43925368ffc511c870 Mon Sep 17 00:00:00 2001 From: Ouroboros Date: Mon, 14 Sep 2026 02:22:51 +0300 Subject: [PATCH] llm_stream: judge form after terminal framing on both accumulators; code-less SSE error is a provider verdict Review round 1 (triad sol/opus/grok, scope astra, adversaries A/B) on bdb95bd5b: - AnthropicAccumulator judges once after message_stop like the Chat accumulator: data after the terminal is noted, non-JSON tool input is left for the terminal verdict, an unusable body raises RejectedProviderStream with its usage (settled, provider_error, no continuation); max_tokens beside a tool block returns like the non-stream path. - ChatAccumulator: a choice-set mismatch after [DONE] is a rejection, not an unknown outcome (a lone foreign index on an n=1 stream is remapped and disclosed); a clean close after every choice finished is terminal framing too (disclosed); a non-text scalar after text keeps the first shape; a tool call without `type` is structurally complete (id, name, arguments). - ProviderStreamError without an HTTP-shaped code is `stream_rejected` (provider_error, no same-request repeat, never unknown); the chat error frame's own usage snapshot settles the attempt; an aborted native stream keeps its unresolved bound. - corpus: .gitattributes keeps recorded bytes exact (-text); parity also compares per-record reasoning_details key sets; README documents each case's request knob; pins for the finish_reason conflict and the later usage snapshot; adjacency docstring states that discreteness of opaque payloads rides on the provider id. Co-authored-by: Ouroboros <311266734+ouroboros-agent@users.noreply.github.com> --- .gitattributes | 6 +- ouroboros/llm_stream.py | 83 +++++++--- ouroboros/usage_accounting.py | 5 +- tests/fixtures/llm_wire/README.md | 21 ++- tests/test_llm_wire_corpus.py | 4 + tests/test_transport_b_stream_deadlines.py | 176 +++++++++++++++++++-- 6 files changed, 248 insertions(+), 47 deletions(-) diff --git a/.gitattributes b/.gitattributes index ab1bc0e8e..8b378bc78 100644 --- a/.gitattributes +++ b/.gitattributes @@ -13,6 +13,6 @@ requirements-runtime.lock text eol=lf # Golden scanner slices are SHA-pinned and must keep identical bytes on Windows. tests/fixtures/skill_publish_scanner/*.fixture text eol=lf -# Recorded provider wire (SSE / JSON bodies) is byte-exact evidence with LF framing; -# keep it LF on every checkout so the replay corpus never sees converted bytes. -tests/fixtures/llm_wire/** text eol=lf +# Recorded provider wire (SSE / JSON bodies) is byte-exact evidence: no line-ending +# conversion at commit or checkout, so the replay corpus always sees the recorded bytes. +tests/fixtures/llm_wire/** -text diff --git a/ouroboros/llm_stream.py b/ouroboros/llm_stream.py index ccc29e337..3bcd0683e 100644 --- a/ouroboros/llm_stream.py +++ b/ouroboros/llm_stream.py @@ -51,7 +51,7 @@ class ProviderStreamError(IncompleteProviderStream): A provider fact, not an unknown outcome: no ``code``; the body's numeric ``error.code`` becomes ``status_code`` so the classifier files it through - its ordinary status ladder. + its ordinary status ladder, and a frame without one is ``stream_rejected``. """ code = "" @@ -61,9 +61,13 @@ class ProviderStreamError(IncompleteProviderStream): error = body.get("error") error = error if isinstance(error, dict) else {} code = error.get("code") - self.status_code = ( - code if isinstance(code, int) and not isinstance(code, bool) and 100 <= code <= 599 else 200 - ) + numeric = isinstance(code, int) and not isinstance(code, bool) and 100 <= code <= 599 + self.status_code = code if numeric else 200 + # Without an HTTP-shaped code the frame is the provider's own terminal + # verdict on this stream (an overload/api_error shape, a finish_reason of + # "error"): filed like a rejected body — provider_error, no same-request + # repeat, the cross-model chain eligible — never an unknown outcome. + self.stream_rejected = not numeric self.type = str(error.get("type") or error.get("code") or "provider_stream_error") self.provider_message = str(error.get("message") or "") self.stream_usage = copy.deepcopy(usage) if isinstance(usage, dict) and usage else None @@ -100,7 +104,13 @@ def _index(value: Any) -> int: def _adjacent(last: Any, item: dict) -> bool: - """Continuation of the last record: ``type`` and ``id`` absent on either side or equal.""" + """Continuation of the last record: ``type`` and ``id`` absent on either side or equal. + + Discreteness of opaque payloads (``reasoning.encrypted`` data, signatures) + rides on the provider's ``id``: two id-less records of one type would fuse. + Every recorded encrypted record carries an id; text records must stay + id-less-mergeable, which is the #856 fix itself. + """ return isinstance(last, dict) and all( last.get(key) is None or item.get(key) is None or last.get(key) == item.get(key) for key in ("type", "id")) @@ -146,8 +156,8 @@ def _delta(target: dict, update: dict, note: Callable[[str], None]) -> None: target[key] = copy.deepcopy(value) elif current != value: note(f"{key}: {current!r} then {value!r}; kept first") - elif isinstance(current, (dict, list)): - note(f"{key}: scalar delta onto {type(current).__name__}; kept first shape") + elif isinstance(current, (dict, list, str)): + note(f"{key}: {type(value).__name__} delta onto {type(current).__name__}; kept first shape") else: target[key] = copy.deepcopy(value) @@ -192,7 +202,9 @@ def _tool_call_problem(call: Any) -> str: return f"is {type(call).__name__}, not an object" if not isinstance(call.get("id"), str) or not call["id"]: return "id missing" - kind = call.get("type") + # Structural completeness is id + name + arguments; a fragment set that never + # carried ``type`` is the same shape the non-stream path executes. + kind = call.get("type") or ("function" if isinstance(call.get("function"), dict) else None) if kind not in {"function", "custom"}: return f"type {kind!r}" payload = call.get(kind) @@ -246,7 +258,7 @@ class ChatAccumulator(_Accumulator): if isinstance(chunk.get("error"), dict): _snapshot(self.body, {key: value for key, value in chunk.items() if key not in {"choices", "usage", "object"}}) - raise ProviderStreamError(chunk, usage=self.body.get("usage")) + raise ProviderStreamError(chunk, usage=chunk.get("usage") or self.body.get("usage")) for key, value in chunk.items(): if key in {"choices", "object", "obfuscation"} or value is None: continue @@ -268,7 +280,7 @@ class ChatAccumulator(_Accumulator): self._note(f"choice at position {position} is not an object; skipped") return index = update.get("index") - if not _valid_index(index): + if not _valid_index(index) or (self.expected_choices == 1 and index != 0): index = 0 if self.expected_choices == 1 else position self._note(f"choice at position {position}: index {update.get('index')!r}; used {index}") delta = update.get("delta") @@ -305,8 +317,16 @@ class ChatAccumulator(_Accumulator): def result(self) -> dict: """Judge once, after terminal framing: unknown outcome vs. unusable body.""" - if not self.done or set(self.choices) != set(range(self.expected_choices)): - raise IncompleteProviderStream("Stream ended without complete terminal framing") + if not self.done: + # A body the provider closed cleanly after every choice finished is + # terminal framing too ([DONE] is the other witness); a close before + # that is the one unknown outcome this assembler still reports. + if not self.choices or any(not choice.get("finish_reason") for choice in self.choices.values()): + raise IncompleteProviderStream("Stream ended without complete terminal framing") + self._note("stream closed without [DONE] after every choice finished") + self.done = True + if set(self.choices) != set(range(self.expected_choices)): + self._reject(f"choices {sorted(self.choices)} present, {self.expected_choices} expected") body = self.partial() for choice in body["choices"]: path = f"choice {choice['index']}" @@ -354,7 +374,8 @@ class AnthropicAccumulator(_Accumulator): def accept(self, event: str, data: str) -> None: if self.done: - raise IncompleteProviderStream("Data after native stream terminal") + self._note("data after message_stop; ignored") + return chunk = json.loads(data) if not isinstance(chunk, dict): raise IncompleteProviderStream("Native stream chunk is not an object") @@ -362,6 +383,9 @@ class AnthropicAccumulator(_Accumulator): if event and event != kind: raise IncompleteProviderStream("Native SSE event differs from payload type") if kind == "error": + # message_start's usage is a lower-bound snapshot (final output tokens + # arrive in message_delta): an aborted native stream keeps its + # unresolved upper bound rather than settling on an understatement. raise ProviderStreamError(chunk) if kind == "ping": return @@ -386,9 +410,11 @@ class AnthropicAccumulator(_Accumulator): if kind == "content_block_stop": self.open_blocks.remove(index) if index in self.inputs: - block["input"] = json.loads(self.inputs[index]) - if not isinstance(block["input"], dict): - raise IncompleteProviderStream("Native tool input is not an object") + try: + block["input"] = json.loads(self.inputs[index]) + except ValueError: + self._note(f"block {index}: tool input is not JSON; left for the terminal verdict") + block["input"] = self.inputs[index] return delta = chunk.get("delta") if not isinstance(delta, dict): @@ -415,24 +441,33 @@ class AnthropicAccumulator(_Accumulator): # Future non-content events are retained in the exact wire evidence. def result(self) -> dict: - if (not self.done or not self.body or self.open_blocks or not self.body.get("stop_reason") - or set(self.blocks) != set(range(len(self.blocks)))): + """Judge once, after ``message_stop``: unknown outcome vs. unusable body.""" + if not self.done: raise IncompleteProviderStream("Native stream ended without complete terminal framing") - if self.body.get("stop_reason") == "max_tokens" and any( - block.get("type") == "tool_use" for block in self.blocks.values()): - raise IncompleteProviderStream("Native output exhausted while producing tool calls") - for block in self.blocks.values(): + if not self.body: + self._reject("message_stop without a message_start") + if self.open_blocks: + self._reject(f"blocks {sorted(self.open_blocks)} never stopped") + if not self.body.get("stop_reason"): + self._reject("no stop_reason after message_stop") + if set(self.blocks) != set(range(len(self.blocks))): + self._reject(f"block indices {sorted(self.blocks)} are not contiguous") + for index, block in self.blocks.items(): kind = block.get("type") if kind in {"tool_use", "server_tool_use"} and ( not isinstance(block.get("id"), str) or not block["id"] or not isinstance(block.get("name"), str) or not block["name"] or not isinstance(block.get("input"), dict)): - raise IncompleteProviderStream("Incomplete native tool block") + self._reject(f"block {index}: tool block lacks id, name or an object input") if kind == "thinking" and (not isinstance(block.get("thinking"), str) or not isinstance(block.get("signature"), str) or not block["signature"]): - raise IncompleteProviderStream("Native thinking block lacks its complete signature") + self._reject(f"block {index}: thinking block lacks its complete signature") return self.partial() + def _reject(self, detail: str) -> None: + raise RejectedProviderStream(f"Native stream rejected after message_stop: {detail}", + usage=self.body.get("usage"), anomalies=self.anomalies) + def partial(self) -> dict: return {**copy.deepcopy(self.body), "content": [copy.deepcopy(self.blocks[key]) for key in sorted(self.blocks)]} diff --git a/ouroboros/usage_accounting.py b/ouroboros/usage_accounting.py index 6f1fca74c..ddfb5a5fe 100644 --- a/ouroboros/usage_accounting.py +++ b/ouroboros/usage_accounting.py @@ -1338,8 +1338,9 @@ def _terminalize_failed_attempt(reservation: AttemptReservation, exc: BaseExcept return "settled" elif isinstance(stream_usage, dict) and stream_usage: # The usage frame was read before the body was judged unusable: money is - # known, so settle exactly as a successful response does (same extractor, - # same cost derivation); only the answer is missing. + # known, so settle through the success path's extractor and cost derivation + # from that frame alone (no assembled-body facts such as service_tier, no + # cache-TTL injection or token-density observation); only the answer is missing. usage, cost, final = usage_from_response({"usage": stream_usage}) settle_attempt(reservation, dict(usage or {}), cost_usd=cost, cost_final=final) return "settled" diff --git a/tests/fixtures/llm_wire/README.md b/tests/fixtures/llm_wire/README.md index 78aa29b4f..718b78a24 100644 --- a/tests/fixtures/llm_wire/README.md +++ b/tests/fixtures/llm_wire/README.md @@ -27,5 +27,22 @@ spacing on frames the redaction did not touch); a redacted `data:` line is re-se compactly. Non-stream bodies keep their leading keep-alive whitespace. Fixture pairs are separate live requests, so parity is structural (message keys, tool -names, argument key sets, `reasoning_details` type sequence, finish reason, usage keys), -never exact text. +names, argument key sets, `reasoning_details` type sequence and per-record key sets, +finish reason, usage keys), never exact text. + +Cases — the same synthetic prompt and `lookup` tool everywhere; what differs is the request +knob (or simply the reply the model chose to give on that run): + +| case | route / model | request knob | what the reply exercises | +|---|---|---|---| +| `tool_stream` (+ `tool_nonstream`) | openrouter / gemini-3.8-flash | `reasoning: {"enabled": true}` | one call, a single `reasoning.encrypted` record | +| `tool_stream_longreasoning` | openrouter / gemini-3.8-flash | `reasoning: {"effort": "low"}` | one call; visible `reasoning.text` then `reasoning.encrypted` — the #856 minimal reproducer (type transition inside one `index`) | +| `multicall_stream` (+ `multicall_nonstream`) | openrouter / gemini-3.8-flash | `reasoning: {"effort": "high"}` | two calls in one turn | +| `secondturn_stream` | openrouter / gemini-3.8-flash | `reasoning: {"effort": "high"}`, second turn (assistant reply + `lookup` result replayed) | continuation with replayed reasoning | +| `tool_stream` | openrouter / grok-4.6 | `reasoning: {"effort": "low"}` | `reasoning.summary` then `reasoning.encrypted` (`rs_…`) | +| `tool_stream`, `multi_stream` | openrouter / gpt-5.6-sol | `reasoning: {"effort": "low"}` / `{"effort": "medium"}` | one / two calls, one `reasoning.encrypted` | +| `tool_stream` (+ `tool_nonstream`) | openrouter / claude-sonnet-5 | `reasoning: {"effort": "high"}` | two calls, `reasoning.text` with a signature | +| `refusal_stream`, `refusal_stream_v2` | openrouter / claude-fable-5 | `reasoning: {"effort": "high"}` | a refusal: `finish_reason: content_filter`, no reasoning | +| `custom_stream` (+ `custom_nonstream`) | openai / gpt-5.6-terra | `tools[0].type = "custom"` with the runtime's `_CUSTOM_FORMAT` | a custom-tool call | +| `function_none_stream` | openai / gpt-5.6-terra | `reasoning_effort: "none"` | two function calls, no reasoning | +| `native_stream` (+ `native_nonstream`) | anthropic / claude-sonnet-5 (Messages) | `thinking: {"type": "adaptive"}` | thinking + text + two `tool_use` blocks | diff --git a/tests/test_llm_wire_corpus.py b/tests/test_llm_wire_corpus.py index 7148908fa..8d307b894 100644 --- a/tests/test_llm_wire_corpus.py +++ b/tests/test_llm_wire_corpus.py @@ -181,6 +181,10 @@ def test_recorded_stream_matches_its_nonstream_sibling_structurally(stream, sibl assert [set(_call_args(call)) for call in streamed_calls] == [set(_call_args(call)) for call in json_calls] assert ([detail["type"] for detail in streamed_msg.get("reasoning_details") or []] == [detail["type"] for detail in json_msg.get("reasoning_details") or []]) + # Per-record key sets: a streamed record that lost its opaque continuation payload + # (``signature``/``data``/``id``) would differ from the non-stream sibling here. + assert ([sorted(detail) for detail in streamed_msg.get("reasoning_details") or []] + == [sorted(detail) for detail in json_msg.get("reasoning_details") or []]) assert streamed_usage["response_finish_reason"] == json_usage["response_finish_reason"] assert set(streamed_usage) - {"stream_receipt"} == set(json_usage) diff --git a/tests/test_transport_b_stream_deadlines.py b/tests/test_transport_b_stream_deadlines.py index 46cb7b502..ba6bc1abf 100644 --- a/tests/test_transport_b_stream_deadlines.py +++ b/tests/test_transport_b_stream_deadlines.py @@ -298,20 +298,59 @@ def test_terminal_length_with_a_partial_tool_call_returns_like_non_stream(isolat assert [row["state"] for row in rows(isolated)] == ["reserved", "dispatched", "settled"] +def test_clean_close_after_every_choice_finished_is_terminal_framing(isolated): + """The provider closed the body cleanly after the choice's ``finish_reason`` but never sent + ``[DONE]``: both witnesses of terminal framing count, so the reply is complete (disclosed in the + receipt), settles, and is not an unknown outcome.""" + wire = sse(chunk({"role": "assistant", "content": "done"}, "stop", usage=completion()["usage"]), done=False) + result = run_driver(lambda **kw: WireResponse(wire), payload(stream=True), target()).model_dump() + msg, usage = LLMClient()._normalize_remote_response(result, target(), skip_cost_fetch=True) + assert msg["content"] == "done" and result["_stream_receipt"]["complete"] is True + assert usage["stream_receipt"]["anomalies"] == {"count": 1, "first": [ + "stream closed without [DONE] after every choice finished"]} + assert rows(isolated)[-1]["state"] == "settled" + + +def test_tool_call_without_type_is_structurally_complete(isolated): + """``type`` is not part of structural completeness (id, name, arguments are): a call whose + fragments never carried it returns exactly as the non-stream path returns it.""" + wire = sse(chunk({"tool_calls": [{"index": 0, "id": "t", "function": {"name": "lookup", "arguments": '{"q":"x"}'}}]}, + "tool_calls", usage=completion()["usage"])) + result = run_driver(lambda **kw: WireResponse(wire), payload(stream=True), target()).model_dump() + assert result["choices"][0]["message"]["tool_calls"] == [ + {"id": "t", "function": {"name": "lookup", "arguments": '{"q":"x"}'}}] + assert result["_stream_receipt"]["anomalies"]["count"] == 0 and rows(isolated)[-1]["state"] == "settled" + + +def test_later_usage_snapshot_overrides_an_earlier_one(isolated): + """Usage frames are cumulative snapshots: the last one read is the one settled, so an early + partial snapshot never becomes the attempt's cost.""" + early = {"id": "gen-test", "object": "chat.completion.chunk", "model": "vendor/test-stream", "choices": [], + "usage": {"prompt_tokens": 10, "completion_tokens": 0, "cost": 0.01}} + wire = sse(early, chunk({"role": "assistant", "content": "done"}), chunk({}, "stop", usage=completion()["usage"])) + result = run_driver(lambda **kw: WireResponse(wire), payload(stream=True), target()).model_dump() + assert result["usage"] == completion()["usage"] + assert rows(isolated)[-1]["state"] == "settled" and rows(isolated)[-1]["cost_usd"] == 0.25 + + def test_identity_conflict_is_forgiven_and_disclosed_in_the_receipt(isolated): - """First value wins for identity scalars (a tool call's ``id``, the envelope ``id``); - each conflict is a fact in ``usage["stream_receipt"]["anomalies"]``, never a raise.""" + """First value wins for identity scalars (a tool call's ``id``, the envelope ``id``) and for a + choice's ``finish_reason``; each conflict is a fact in ``usage["stream_receipt"]["anomalies"]``, + never a raise.""" wire = sse(chunk({"tool_calls": [{"index": 0, "id": "call_a", "type": "function", "function": {"name": "lookup", "arguments": '{"q":'}}]}), chunk({"tool_calls": [{"index": 0, "id": "call_b", "function": {"arguments": '"ok"}'}}]}, - "tool_calls", usage=completion()["usage"], id="gen-other")) + "tool_calls", usage=completion()["usage"], id="gen-other"), + chunk({}, "stop")) result = run_driver(lambda **kw: WireResponse(wire), payload(stream=True), target()).model_dump() msg, usage = LLMClient()._normalize_remote_response(result, target(), skip_cost_fetch=True) assert msg["tool_calls"][0]["id"] == "call_a" and msg["response_id"] == "gen-test" assert json.loads(msg["tool_calls"][0]["function"]["arguments"]) == {"q": "ok"} - assert usage["stream_receipt"]["anomalies"] == {"count": 2, "first": [ + assert result["choices"][0]["finish_reason"] == "tool_calls" + assert usage["stream_receipt"]["anomalies"] == {"count": 3, "first": [ "id: 'gen-test' then 'gen-other'; kept first", "id: 'call_a' then 'call_b'; kept first", + "choice 0: finish_reason 'tool_calls' then 'stop'; kept first", ]} assert rows(isolated)[-1]["state"] == "settled" @@ -332,6 +371,39 @@ def test_finish_frame_without_delta_still_completes_the_choice(isolated): assert rows(isolated)[-1]["state"] == "settled" +def test_lone_choice_with_a_foreign_index_is_remapped_on_a_single_choice_stream(isolated): + """A single-choice reply whose only choice carries index 1 is a form irregularity after + ``[DONE]``: it is remapped to choice 0 and disclosed, never an unknown outcome.""" + wire = sse(chunk({"role": "assistant", "content": "done"}, "stop", index=1, usage=completion()["usage"])) + result = run_driver(lambda **kw: WireResponse(wire), payload(stream=True), target()).model_dump() + assert result["choices"][0]["index"] == 0 and result["choices"][0]["message"]["content"] == "done" + msg, usage = LLMClient()._normalize_remote_response(result, target(), skip_cost_fetch=True) + assert usage["stream_receipt"]["anomalies"]["first"] == ["choice at position 0: index 1; used 0"] + assert rows(isolated)[-1]["state"] == "settled" + + +def test_missing_choice_after_done_is_rejected_not_unknown(isolated): + """``n=2`` with only one choice at ``[DONE]``: the wire is complete, the body is unusable — + a rejection with usage (settled), not an unknown outcome.""" + wire = sse(chunk({"role": "assistant", "content": "only one"}, "stop", usage=completion()["usage"])) + with pytest.raises(RejectedProviderStream) as caught: + run_driver(lambda **kw: WireResponse(wire), payload(stream=True, n=2), target()) + assert "choices [0] present, 2 expected" in str(caught.value) and caught.value.stream_receipt["complete"] is True + assert rows(isolated)[-1]["state"] == "settled" + + +def test_non_text_scalar_after_text_keeps_the_first_shape(isolated): + """First shape wins in both directions: text established by earlier frames is not + overwritten by a later malformed non-text scalar; the conflict is disclosed.""" + wire = sse(chunk({"role": "assistant", "content": "ok"}), chunk({"content": 7}), + chunk({}, "stop", usage=completion()["usage"])) + result = run_driver(lambda **kw: WireResponse(wire), payload(stream=True), target()).model_dump() + msg, usage = LLMClient()._normalize_remote_response(result, target(), skip_cost_fetch=True) + assert msg["content"] == "ok" + assert usage["stream_receipt"]["anomalies"] == {"count": 1, "first": ["content: int delta onto str; kept first shape"]} + assert rows(isolated)[-1]["state"] == "settled" + + def test_index_less_tool_call_fragment_continues_the_last_call(isolated): """A tool-call delta without ``index`` is a fragment of the last call when its type/id are compatible (merged, disclosed); a fresh id is a new call (appended, disclosed).""" @@ -393,6 +465,43 @@ def test_mid_stream_error_chunk_classifies_by_body_and_settles_only_with_usage(i "provider_transient", True, 502) +@pytest.mark.parametrize("usage_first", [False, True]) +@pytest.mark.parametrize("shape", ["finish_reason_error", "code_less_error_frame"]) +def test_code_less_stream_error_is_a_provider_verdict_not_an_unknown_outcome(isolated, usage_first, shape): + """An SSE error without an HTTP-shaped code — the provider's own ``finish_reason: "error"``, or an + ``{"error": {"type": ...}}`` frame (the overload/api_error shape) — is the provider's terminal verdict on + this stream: ``provider_error`` with no same-request repeat, whether or not a usage frame was read + first; it never reopens the unknown-outcome continuation. A usage snapshot carried by the error frame + itself settles the attempt.""" + frames = [chunk({"content": "par"})] + if usage_first: + frames.append({"id": "gen-test", "choices": [], "usage": completion()["usage"]}) + if shape == "finish_reason_error": + frames.append(chunk({"content": ""}, "error")) + else: + frames.append({"id": "gen-test", "error": {"type": "overloaded_error", "message": "Overloaded"}, + **({"usage": completion()["usage"]} if usage_first else {})}) + with pytest.raises(ProviderStreamError) as caught: + run_driver(lambda **kw: WireResponse(sse(*frames, done=False)), payload(stream=True), target()) + exc = caught.value + assert exc.stream_rejected and exc.code == "" and exc.status_code == 200 + assert rows(isolated)[-1]["state"] == ("settled" if usage_first else "unresolved") + classification = classify_llm_exception(exc) + assert (classification.kind, classification.retry_same_request) == ("provider_error", False) + + +def test_error_frame_carrying_its_own_usage_settles_the_attempt(isolated): + """The usage snapshot on the error frame itself is the latest money fact: it settles the attempt + instead of leaving an unresolved upper bound behind a known outcome.""" + frames = [chunk({"content": "par"}), + {"id": "gen-test", "error": {"code": 502, "message": "Upstream provider error"}, + "usage": completion()["usage"]}] + with pytest.raises(ProviderStreamError) as caught: + run_driver(lambda **kw: WireResponse(sse(*frames, done=False)), payload(stream=True), target()) + assert caught.value.stream_usage == completion()["usage"] + assert rows(isolated)[-1]["state"] == "settled" and rows(isolated)[-1]["cost_usd"] == 0.25 + + @pytest.mark.parametrize("asynchronous", [False, True]) @pytest.mark.parametrize("field", ["stream", "stream_options"]) def test_stream_rejection_uses_existing_wire_recovery(isolated, asynchronous, field): @@ -624,17 +733,52 @@ def test_native_incomplete_blocks_or_message_cannot_return_tools(isolated, monke assert rows(isolated)[-1]["state"] == "unresolved" +def test_native_unusable_body_after_message_stop_is_rejected_and_settled(isolated, monkeypatch): + """A thinking block that never received its signature is unusable, but ``message_stop`` and the + usage arrived: the native path judges once after terminal framing exactly like the Chat path — + ``RejectedProviderStream`` with the usage, a settled ledger row, and ``provider_error`` (no retry of + the same request, no unknown-outcome continuation).""" + import requests + events, _expected = native_events() + events = [(kind, body) for kind, body in events + if not (kind == "content_block_delta" and (body.get("delta") or {}).get("type") == "signature_delta")] + monkeypatch.setattr(requests, "post", lambda *a, **k: WireResponse(sse(*events, done=False))) + with pytest.raises(RejectedProviderStream) as caught: + LLMClient()._chat_anthropic(target("anthropic"), MESSAGES, TOOLS, "high", 1024, "auto", stream=True) + assert "thinking block lacks its complete signature" in str(caught.value) + assert caught.value.stream_usage["input_tokens"] == 5 and not hasattr(caught.value, "code") + assert rows(isolated)[-1]["state"] == "settled" + verdict = classify_llm_exception(caught.value) + assert (verdict.kind, verdict.retry_same_request) == ("provider_error", False) + + +def test_native_max_tokens_with_tool_blocks_returns_like_non_stream(isolated, monkeypatch): + """``stop_reason == max_tokens`` beside a complete tool block is a finished reply the loop reads + (parity with the non-stream native path, which surfaces ``stop_reason``); it is not a rejection.""" + import requests + events, expected = native_events() + for kind, body in events: + if kind == "message_delta": + body["delta"]["stop_reason"] = "max_tokens" + monkeypatch.setattr(requests, "post", lambda *a, **k: WireResponse(sse(*events, done=False))) + message, usage = LLMClient()._chat_anthropic(target("anthropic"), MESSAGES, TOOLS, "high", 1024, "auto", stream=True) + assert message["stop_reason"] == "max_tokens" and message["tool_calls"][0]["function"]["name"] == "lookup" + assert rows(isolated)[-1]["state"] == "settled" + + @pytest.mark.parametrize("before_content", [False, True]) @pytest.mark.parametrize("error,status,kind", [ - ({"type": "overloaded_error", "message": "Overloaded"}, 200, "provider_outcome_unknown"), + ({"type": "overloaded_error", "message": "Overloaded"}, 200, "provider_error"), ({"type": "api_error", "code": 502, "message": "Bad gateway"}, 502, "provider_transient"), ]) def test_native_http_200_sse_error_keeps_producer_facts_and_custody(isolated, monkeypatch, before_content, error, status, kind): """An explicit SSE error is a provider fact, not ``model_outcome_unknown``: a numeric body code becomes the status the classifier files (502 → provider_transient); a - type-only body keeps status 200, so unresolved custody alone still reads unknown. - Custody stays unresolved either way because no usage frame was read.""" + type-only body (the overload shape) is the provider's own verdict on the stream + (``stream_rejected`` → provider_error, no same-request repeat), never an unknown + outcome. Custody stays unresolved either way: no final usage frame was read, and + ``message_start``'s snapshot is only a lower bound.""" import requests events, _ = native_events() event = {"type": "error", "error": error} @@ -645,6 +789,7 @@ def test_native_http_200_sse_error_keeps_producer_facts_and_custody(isolated, mo exc = caught.value assert exc.body == event and exc.type == error["type"] assert exc.code == "" and exc.status_code == status and exc.stream_usage is None + assert exc.stream_rejected is (status == 200) assert exc.stream_receipt["generation_id"] == "header-generation" and exc.stream_receipt["complete"] is False assert response.closed and rows(isolated)[-1]["state"] == "unresolved" assert classify_llm_exception(exc).kind == kind @@ -799,16 +944,15 @@ def test_actual_sdk_loopback_sse_and_cleanup(isolated, monkeypatch, asynchronous operation = lambda: asyncio.run(call_and_close()) else: operation = lambda: client.chat(**kwargs) - if terminal: - msg, usage = operation() - assert msg["content"] == "done" - assert usage["stream_receipt"]["generation_id"] == "loopback-generation" - else: - with pytest.raises(IncompleteProviderStream) as caught: - operation() - assert caught.value.stream_receipt["generation_id"] == "loopback-generation" + msg, usage = operation() + assert msg["content"] == "done" + assert usage["stream_receipt"]["generation_id"] == "loopback-generation" + # A body closed cleanly after every choice finished is terminal framing even without + # ``[DONE]`` (a compatible endpoint that omits it must not turn every reply into an + # unknown outcome); the missing witness is disclosed in the receipt. + assert usage["stream_receipt"]["anomalies"]["count"] == (0 if terminal else 1) assert observed[0]["stream"] is True and observed[0]["stream_options"]["include_usage"] is True - assert [row["state"] for row in rows(isolated)] == ["reserved", "dispatched", "settled" if terminal else "unresolved"] + assert [row["state"] for row in rows(isolated)] == ["reserved", "dispatched", "settled"] finally: for sdk in client._remote_clients.values(): sdk.close()