From 19b3501fc25bbaa107103365bcfb5ea94ddce964 Mon Sep 17 00:00:00 2001 From: Ben Clarke Date: Tue, 25 Aug 2026 11:23:04 +0100 Subject: [PATCH] fix: join streamed reasoning into one thought part per block Both stream finalizers passed the reasoning buffer through unjoined, so the aggregated response held one thought part per delta while text was joined into one. `_aggregate_streaming_thought_parts` already does the joining and splits on `thought_signature`, so Anthropic block boundaries survive. --- src/google/adk/models/lite_llm.py | 12 ++- tests/unittests/models/test_litellm.py | 107 +++++++++++++++++++++++++ 2 files changed, 117 insertions(+), 2 deletions(-) diff --git a/src/google/adk/models/lite_llm.py b/src/google/adk/models/lite_llm.py index 85e90d4e6a..d4d6c76763 100644 --- a/src/google/adk/models/lite_llm.py +++ b/src/google/adk/models/lite_llm.py @@ -3232,7 +3232,11 @@ def _finalize_tool_call_response( tool_calls=tool_calls, ), model_version=model_version, - thought_parts=list(reasoning_parts) if reasoning_parts else None, + thought_parts=( + _aggregate_streaming_thought_parts(reasoning_parts) + if reasoning_parts + else None + ), ) mapped_finish_reason = _map_finish_reason(finish_reason) llm_response.finish_reason = mapped_finish_reason @@ -3256,7 +3260,11 @@ def _finalize_text_response( content=message_content, ), model_version=model_version, - thought_parts=list(reasoning_parts) if reasoning_parts else None, + thought_parts=( + _aggregate_streaming_thought_parts(reasoning_parts) + if reasoning_parts + else None + ), ) mapped_finish_reason = _map_finish_reason(finish_reason) llm_response.finish_reason = mapped_finish_reason diff --git a/tests/unittests/models/test_litellm.py b/tests/unittests/models/test_litellm.py index 6c89f2a725..60685e0c72 100644 --- a/tests/unittests/models/test_litellm.py +++ b/tests/unittests/models/test_litellm.py @@ -6856,6 +6856,113 @@ async def test_streaming_text_assembled_from_many_fragments( assert aggregated[0].content.parts[0].text == full_text +def _reasoning_stream_chunks(deltas, finish_reason="stop"): + stream = [ + ModelResponseStream( + choices=[ + StreamingChoices( + finish_reason=None, delta=Delta(role="assistant", **kwargs) + ) + ] + ) + for kwargs in deltas + ] + stream.append( + ModelResponseStream( + choices=[StreamingChoices(finish_reason=finish_reason, delta=Delta())] + ) + ) + return stream + + +@pytest.mark.asyncio +async def test_streaming_reasoning_assembled_from_many_fragments( + mock_completion, lite_llm_instance +): + # Providers that carry no thought signature (xAI, OpenAI, Ollama) stream + # reasoning one delta per token. The aggregated response has to join them the + # way the text buffer does, or the stored event holds one part per token. + full_reasoning = "".join(f"token-{i} " for i in range(50)) + fragments = _split_into_chunks( + full_reasoning, [7] * (len(full_reasoning) // 7) + ) + mock_completion.return_value = iter( + _reasoning_stream_chunks( + [{"reasoning_content": fragment} for fragment in fragments] + ) + ) + + responses = [ + r + async for r in lite_llm_instance.generate_content_async( + LLM_REQUEST_WITH_FUNCTION_DECLARATION, stream=True + ) + ] + + partials = [r for r in responses if r.partial] + aggregated = [r for r in responses if not r.partial] + assert [p.content.parts[0].text for p in partials] == fragments + assert len(aggregated) == 1 + parts = aggregated[0].content.parts + assert len(parts) == 1 + assert parts[0].thought is True + assert parts[0].text == full_reasoning + + +@pytest.mark.asyncio +async def test_streaming_reasoning_keeps_thinking_block_boundaries( + mock_completion, lite_llm_instance +): + # Anthropic closes each thinking block with a signature-only delta, and the + # signature only matches the text of its own block, so blocks must aggregate + # one part each rather than collapsing into one. + mock_completion.return_value = iter( + _reasoning_stream_chunks([ + { + "thinking_blocks": [ + {"type": "thinking", "thinking": "First half "} + ] + }, + { + "thinking_blocks": [ + {"type": "thinking", "thinking": "of block one."} + ] + }, + { + "thinking_blocks": [{ + "type": "thinking", + "thinking": "", + "signature": "c2lnLW9uZQ==", + }] + }, + {"thinking_blocks": [{"type": "thinking", "thinking": "Block two."}]}, + { + "thinking_blocks": [{ + "type": "thinking", + "thinking": "", + "signature": "c2lnLXR3bw==", + }] + }, + ]) + ) + + responses = [ + r + async for r in lite_llm_instance.generate_content_async( + LLM_REQUEST_WITH_FUNCTION_DECLARATION, stream=True + ) + ] + + aggregated = [r for r in responses if not r.partial] + assert len(aggregated) == 1 + parts = aggregated[0].content.parts + assert len(parts) == 2 + assert parts[0].text == "First half of block one." + assert parts[0].thought_signature == b"sig-one" + assert parts[1].text == "Block two." + assert parts[1].thought_signature == b"sig-two" + + @pytest.mark.asyncio async def test_streaming_buffers_hold_fragments_instead_of_growing_copies( mock_completion, lite_llm_instance