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>
This commit is contained in:
Ouroboros 2026-09-14 02:22:51 +03:00
parent 37f25fb63f
commit 3ce1d96f4d
6 changed files with 248 additions and 47 deletions

6
.gitattributes vendored
View file

@ -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

View file

@ -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)]}

View file

@ -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"

View file

@ -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 |

View file

@ -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)

View file

@ -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()