From 602727b178a85026b7c416675763eef32ceabed4 Mon Sep 17 00:00:00 2001 From: Vineeth Sai Date: Mon, 24 Aug 2026 17:24:35 -0700 Subject: [PATCH] fix(streaming): stop stream_chunk_builder duplicating Anthropic thinking text Anthropic re-sends a thinking block's whole text alongside its signature. `ModelResponseIterator._handle_content_block_delta` builds that signature chunk by joining every prior thinking delta, so `get_combined_thinking_content` appended the full text on top of the parts it had already accumulated and the reassembled reasoning came back twice. Skip the append when the block carries a signature and its text is exactly what has been accumulated so far. A provider that sends only the last increment alongside the signature still has that increment appended. Rebased onto litellm_internal_staging: the branch had drifted far enough that it was carrying twelve commits belonging to other PRs, so the diff showed six unrelated e2e files. Rebuilt on the current tip so it is the two files it should have been. --- .../streaming_chunk_builder_utils.py | 7 +- .../test_streaming_chunk_builder_utils.py | 93 +++++++++++++++++++ 2 files changed, 98 insertions(+), 2 deletions(-) diff --git a/litellm/litellm_core_utils/streaming_chunk_builder_utils.py b/litellm/litellm_core_utils/streaming_chunk_builder_utils.py index ee0518c4aec..9a2869e2a40 100644 --- a/litellm/litellm_core_utils/streaming_chunk_builder_utils.py +++ b/litellm/litellm_core_utils/streaming_chunk_builder_utils.py @@ -631,9 +631,12 @@ class ChunkProcessor: ) else: thinking_text = thinking_block.get("thinking", None) - if thinking_text: - current_thinking_text_parts.append(thinking_text) signature = thinking_block.get("signature", None) + already_accumulated = bool(signature) and thinking_text == "".join( + current_thinking_text_parts + ) + if thinking_text and not already_accumulated: + current_thinking_text_parts.append(thinking_text) if signature: current_signature = signature _flush_thinking_block() diff --git a/tests/test_litellm/litellm_core_utils/test_streaming_chunk_builder_utils.py b/tests/test_litellm/litellm_core_utils/test_streaming_chunk_builder_utils.py index 4b5b51cb4b8..5dc273dba69 100644 --- a/tests/test_litellm/litellm_core_utils/test_streaming_chunk_builder_utils.py +++ b/tests/test_litellm/litellm_core_utils/test_streaming_chunk_builder_utils.py @@ -155,6 +155,99 @@ def test_get_combined_tool_content(): ] +def test_get_combined_thinking_content_does_not_duplicate_resent_thinking(): + """Anthropic re-sends a block's whole thinking text alongside its signature. + + ``ModelResponseIterator._handle_content_block_delta`` builds the signature + chunk by joining every prior thinking delta, so appending it on top of the + parts already accumulated emitted the reasoning twice. + """ + base_chunk = { + "id": "chatcmpl-123", + "object": "chat.completion.chunk", + "created": 1234567890, + "model": "claude-sonnet-4-20250514", + } + + def make_chunk(**delta_kwargs): + return ModelResponseStream( + **base_chunk, + choices=[ + StreamingChoices( + index=0, + delta=Delta(**delta_kwargs), + finish_reason=None, + ) + ], + ) + + parts = ["Let me ", "work through ", "this step by step."] + chunks = [ + make_chunk(role="assistant", content=None), + *[make_chunk(thinking_blocks=[{"type": "thinking", "thinking": part}]) for part in parts], + make_chunk( + thinking_blocks=[ + { + "type": "thinking", + "thinking": "".join(parts), + "signature": "sig_block1", + } + ] + ), + ] + thinking_chunks = [ + chunk for chunk in chunks if chunk["choices"][0]["delta"].get("thinking_blocks") + ] + processor = ChunkProcessor(chunks=chunks) + + result = processor.get_combined_thinking_content(thinking_chunks) + + assert result is not None + assert len(result) == 1 + assert result[0]["thinking"] == "".join(parts) + assert result[0]["signature"] == "sig_block1" + + +def test_get_combined_thinking_content_keeps_a_genuine_final_increment(): + """A provider that sends only the last increment with the signature must + still have that increment appended, not swallowed by the de-duplication. + """ + base_chunk = { + "id": "chatcmpl-123", + "object": "chat.completion.chunk", + "created": 1234567890, + "model": "claude-sonnet-4-20250514", + } + + def make_chunk(**delta_kwargs): + return ModelResponseStream( + **base_chunk, + choices=[ + StreamingChoices( + index=0, + delta=Delta(**delta_kwargs), + finish_reason=None, + ) + ], + ) + + chunks = [ + make_chunk(thinking_blocks=[{"type": "thinking", "thinking": "Let me "}]), + make_chunk( + thinking_blocks=[ + {"type": "thinking", "thinking": "think.", "signature": "sig_block1"} + ] + ), + ] + processor = ChunkProcessor(chunks=chunks) + + result = processor.get_combined_thinking_content(chunks) + + assert result is not None + assert len(result) == 1 + assert result[0]["thinking"] == "Let me think." + + def test_get_combined_thinking_content_preserves_interleaved_blocks(): base_chunk = { "id": "chatcmpl-123",