From 2a8e29cc0be05cce13cc47ebaa523030d3259dd3 Mon Sep 17 00:00:00 2001 From: Napuh Date: Sun, 26 Jul 2026 12:30:02 +0200 Subject: [PATCH 1/2] fix(anthropic): split mixed reasoning stream chunks --- .../adapters/streaming_iterator.py | 45 ++++++++++++++- .../test_streaming_iterator_first_delta.py | 56 ++++++++++++++----- 2 files changed, 85 insertions(+), 16 deletions(-) diff --git a/litellm/llms/anthropic/experimental_pass_through/adapters/streaming_iterator.py b/litellm/llms/anthropic/experimental_pass_through/adapters/streaming_iterator.py index f02333c34c86..2103bac5ebf8 100644 --- a/litellm/llms/anthropic/experimental_pass_through/adapters/streaming_iterator.py +++ b/litellm/llms/anthropic/experimental_pass_through/adapters/streaming_iterator.py @@ -96,6 +96,38 @@ def _is_combined(chunk: Any) -> bool: or getattr(delta, "thinking_blocks", None) ) + @staticmethod + def _has_mixed_reasoning_and_text(chunk: Any) -> bool: + choices = getattr(chunk, "choices", None) + if not choices: + return False + delta = getattr(choices[0], "delta", None) + if delta is None: + return False + return bool(getattr(delta, "content", None) and getattr(delta, "reasoning_content", None)) + + @staticmethod + def _clear_usage(chunk: Any) -> None: + if hasattr(chunk, "usage"): + chunk.usage = None + hidden_params = getattr(chunk, "_hidden_params", None) + if isinstance(hidden_params, dict) and "usage" in hidden_params: + chunk._hidden_params = {key: value for key, value in hidden_params.items() if key != "usage"} + + @staticmethod + def _split_mixed_reasoning_and_text(chunk: Any) -> List[Any]: + if not _CombinedChunkSplitter._has_mixed_reasoning_and_text(chunk): + return [chunk] + + reasoning_chunk = copy.deepcopy(chunk) + reasoning_chunk.choices[0].finish_reason = None + reasoning_chunk.choices[0].delta.content = None + _CombinedChunkSplitter._clear_usage(reasoning_chunk) + + text_chunk = copy.deepcopy(chunk) + text_chunk.choices[0].delta.reasoning_content = None + return [reasoning_chunk, text_chunk] + @staticmethod def _split(chunk: Any) -> List[Any]: """Return ``[chunk]``, or ``[content_chunk, finish_chunk]`` if combined.""" @@ -105,6 +137,7 @@ def _split(chunk: Any) -> List[Any]: # Content chunk: keep the delta payload, clear the finish_reason. content_chunk = copy.deepcopy(chunk) content_chunk.choices[0].finish_reason = None + _CombinedChunkSplitter._clear_usage(content_chunk) # Finish chunk: keep finish_reason (and usage), clear the delta payload. finish_chunk = copy.deepcopy(chunk) @@ -127,7 +160,11 @@ def __next__(self) -> Any: if self._sync_iter is None: self._sync_iter = iter(self._stream) chunk = next(self._sync_iter) # propagates StopIteration when exhausted - self._buffer.extend(self._split(chunk)) + self._buffer.extend( + split_chunk + for combined_chunk in self._split(chunk) + for split_chunk in self._split_mixed_reasoning_and_text(combined_chunk) + ) return self._buffer.popleft() def __aiter__(self) -> "AsyncIterator[Any]": @@ -139,7 +176,11 @@ async def __anext__(self) -> Any: if self._async_iter is None: self._async_iter = self._stream.__aiter__() chunk = await self._async_iter.__anext__() # propagates StopAsyncIteration - self._buffer.extend(self._split(chunk)) + self._buffer.extend( + split_chunk + for combined_chunk in self._split(chunk) + for split_chunk in self._split_mixed_reasoning_and_text(combined_chunk) + ) return self._buffer.popleft() diff --git a/tests/test_litellm/llms/anthropic/experimental_pass_through/adapters/test_streaming_iterator_first_delta.py b/tests/test_litellm/llms/anthropic/experimental_pass_through/adapters/test_streaming_iterator_first_delta.py index 19ec1a04b454..3e2ac1ab76e1 100644 --- a/tests/test_litellm/llms/anthropic/experimental_pass_through/adapters/test_streaming_iterator_first_delta.py +++ b/tests/test_litellm/llms/anthropic/experimental_pass_through/adapters/test_streaming_iterator_first_delta.py @@ -65,9 +65,7 @@ def _thinking_chunk(thinking: str, signature: str = "") -> MagicMock: return _make_chunk(Delta(content=None, thinking_blocks=[block])) -def _tool_chunk( - call_id: str, name: Optional[str], arguments: Optional[str] -) -> MagicMock: +def _tool_chunk(call_id: str, name: Optional[str], arguments: Optional[str]) -> MagicMock: return _make_chunk( Delta( content=None, @@ -109,8 +107,7 @@ def _text_deltas(events: List[dict]) -> List[str]: return [ e["delta"]["text"] for e in events - if e.get("type") == "content_block_delta" - and e["delta"].get("type") == "text_delta" + if e.get("type") == "content_block_delta" and e["delta"].get("type") == "text_delta" ] @@ -118,8 +115,7 @@ def _input_json_deltas(events: List[dict]) -> List[str]: return [ e["delta"]["partial_json"] for e in events - if e.get("type") == "content_block_delta" - and e["delta"].get("type") == "input_json_delta" + if e.get("type") == "content_block_delta" and e["delta"].get("type") == "input_json_delta" ] @@ -127,8 +123,7 @@ def _thinking_deltas(events: List[dict]) -> List[str]: return [ e["delta"]["thinking"] for e in events - if e.get("type") == "content_block_delta" - and e["delta"].get("type") == "thinking_delta" + if e.get("type") == "content_block_delta" and e["delta"].get("type") == "thinking_delta" ] @@ -136,8 +131,7 @@ def _signature_deltas(events: List[dict]) -> List[str]: return [ e["delta"]["signature"] for e in events - if e.get("type") == "content_block_delta" - and e["delta"].get("type") == "signature_delta" + if e.get("type") == "content_block_delta" and e["delta"].get("type") == "signature_delta" ] @@ -228,9 +222,7 @@ async def test_first_text_delta_after_tool_use_is_not_dropped_async(): _make_chunk(Delta(content=" Bye.")), _make_chunk(Delta(content=None), finish_reason="stop"), ] - wrapper = AnthropicStreamWrapper( - completion_stream=_AsyncStream(chunks), model="claude-x" - ) + wrapper = AnthropicStreamWrapper(completion_stream=_AsyncStream(chunks), model="claude-x") events = await _drain_async(wrapper) assert _input_json_deltas(events) == ['{"city": "NY"}'] @@ -506,3 +498,39 @@ def test_empty_content_chunk_mid_text_block_is_suppressed_sync(): assert _text_deltas(events) == ["Hi", " there"] _assert_deltas_match_their_block_type(events) + + +def _mixed_reasoning_and_text_chunks() -> List[MagicMock]: + return [ + _make_chunk(Delta(content=None, reasoning_content="First thought.")), + _make_chunk( + Delta(content="Answer.", reasoning_content=" Last thought."), + finish_reason="stop", + ), + ] + + +def _assert_mixed_reasoning_and_text_chunk_is_split(events: List[dict]) -> None: + _assert_deltas_match_their_block_type(events) + assert _thinking_deltas(events) == ["First thought.", " Last thought."] + assert _text_deltas(events) == ["Answer."] + assert [event["type"] for event in events].count("message_delta") == 1 + + +def test_mixed_reasoning_and_text_chunk_is_split_sync(): + wrapper = AnthropicStreamWrapper( + completion_stream=iter(_mixed_reasoning_and_text_chunks()), + model="claude-x", + ) + + _assert_mixed_reasoning_and_text_chunk_is_split(_drain_sync(wrapper)) + + +@pytest.mark.asyncio +async def test_mixed_reasoning_and_text_chunk_is_split_async(): + wrapper = AnthropicStreamWrapper( + completion_stream=_AsyncStream(_mixed_reasoning_and_text_chunks()), + model="claude-x", + ) + + _assert_mixed_reasoning_and_text_chunk_is_split(await _drain_async(wrapper)) From 2b94847364aa284c54ea52ca76721eb330a2c227 Mon Sep 17 00:00:00 2001 From: Napuh Date: Sun, 26 Jul 2026 13:07:49 +0200 Subject: [PATCH 2/2] style: use builtin generic annotation --- .../experimental_pass_through/adapters/streaming_iterator.py | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/litellm/llms/anthropic/experimental_pass_through/adapters/streaming_iterator.py b/litellm/llms/anthropic/experimental_pass_through/adapters/streaming_iterator.py index 2103bac5ebf8..f46f9a617449 100644 --- a/litellm/llms/anthropic/experimental_pass_through/adapters/streaming_iterator.py +++ b/litellm/llms/anthropic/experimental_pass_through/adapters/streaming_iterator.py @@ -115,7 +115,7 @@ def _clear_usage(chunk: Any) -> None: chunk._hidden_params = {key: value for key, value in hidden_params.items() if key != "usage"} @staticmethod - def _split_mixed_reasoning_and_text(chunk: Any) -> List[Any]: + def _split_mixed_reasoning_and_text(chunk: Any) -> list[Any]: if not _CombinedChunkSplitter._has_mixed_reasoning_and_text(chunk): return [chunk]