From 84b8c652e4b90038e28d3233de1e539a4b9ebe8a Mon Sep 17 00:00:00 2001 From: Chesars Date: Sun, 22 Mar 2026 10:33:21 -0300 Subject: [PATCH] fix: preserve tool_use input args in Anthropic adapter streaming When Gemini sends tool call arguments in the same streaming chunk as a content block transition, the Anthropic adapter discarded the processed_chunk containing the input_json_delta. This caused tool_use blocks to arrive with empty input: {}. Queue the processed_chunk alongside the block transition events when it contains input_json_delta data. Applied to both sync and async paths. Fixes #24134 --- .../adapters/streaming_iterator.py | 28 +-- ...al_pass_through_adapters_transformation.py | 171 ++++++++++++++++++ 2 files changed, 187 insertions(+), 12 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 6bddad09f21..806513f36ff 100644 --- a/litellm/llms/anthropic/experimental_pass_through/adapters/streaming_iterator.py +++ b/litellm/llms/anthropic/experimental_pass_through/adapters/streaming_iterator.py @@ -129,8 +129,6 @@ class AnthropicStreamWrapper(AdapterCompletionStreamWrapper): if should_start_new_block and not self.sent_content_block_finish: # Queue the sequence: content_block_stop -> content_block_start - # The trigger chunk itself is not emitted as a delta since the - # content_block_start already carries the relevant information. self.chunk_queue.append( { "type": "content_block_stop", @@ -144,6 +142,14 @@ class AnthropicStreamWrapper(AdapterCompletionStreamWrapper): "content_block": self.current_content_block_start, } ) + # Gemini sends tool args in the same chunk as the block + # transition — queue the delta so it's not lost. + if ( + processed_chunk["type"] == "content_block_delta" + and processed_chunk.get("delta", {}).get("type") == "input_json_delta" + and processed_chunk["delta"].get("partial_json") + ): + self.chunk_queue.append(processed_chunk) self.sent_content_block_finish = False return self.chunk_queue.popleft() @@ -305,18 +311,12 @@ class AnthropicStreamWrapper(AdapterCompletionStreamWrapper): if not self.queued_usage_chunk: if should_start_new_block and not self.sent_content_block_finish: # Queue the sequence: content_block_stop -> content_block_start - # The trigger chunk itself is not emitted as a delta since the - # content_block_start already carries the relevant information. - - # 1. Stop current content block self.chunk_queue.append( { "type": "content_block_stop", "index": max(self.current_content_block_index - 1, 0), } ) - - # 2. Start new content block self.chunk_queue.append( { "type": "content_block_start", @@ -324,11 +324,15 @@ class AnthropicStreamWrapper(AdapterCompletionStreamWrapper): "content_block": self.current_content_block_start, } ) - - # Reset state for new block + # Gemini sends tool args in the same chunk as the block + # transition — queue the delta so it's not lost. + if ( + processed_chunk["type"] == "content_block_delta" + and processed_chunk.get("delta", {}).get("type") == "input_json_delta" + and processed_chunk["delta"].get("partial_json") + ): + self.chunk_queue.append(processed_chunk) self.sent_content_block_finish = False - - # Return the first queued item return self.chunk_queue.popleft() if ( diff --git a/tests/test_litellm/llms/anthropic/experimental_pass_through/adapters/test_anthropic_experimental_pass_through_adapters_transformation.py b/tests/test_litellm/llms/anthropic/experimental_pass_through/adapters/test_anthropic_experimental_pass_through_adapters_transformation.py index ae970e1ff06..3d1a45c6ca1 100644 --- a/tests/test_litellm/llms/anthropic/experimental_pass_through/adapters/test_anthropic_experimental_pass_through_adapters_transformation.py +++ b/tests/test_litellm/llms/anthropic/experimental_pass_through/adapters/test_anthropic_experimental_pass_through_adapters_transformation.py @@ -25,6 +25,7 @@ from litellm.types.utils import ( Function, Message, ModelResponse, + ModelResponseStream, StreamingChoices, Usage, ) @@ -2111,3 +2112,173 @@ class TestTranslateAnthropicOutputFormatToOpenAI: assert self.adapter.translate_anthropic_output_format_to_openai("invalid") is None assert self.adapter.translate_anthropic_output_format_to_openai({"type": "text"}) is None assert self.adapter.translate_anthropic_output_format_to_openai({"type": "json_schema"}) is None + + +class TestAnthropicStreamWrapperToolArgs: + """ + Regression test for https://github.com/BerriAI/litellm/issues/24134 + + When Gemini sends tool call args in the same streaming chunk as a content + block transition, the Anthropic adapter was discarding the processed_chunk + containing input_json_delta. This verifies the args are preserved. + """ + + def _parse_sse_events(self, raw_bytes_list): + """Parse SSE bytes from the stream wrapper into event dicts.""" + import json + + events = [] + for raw in raw_bytes_list: + if not isinstance(raw, bytes): + continue + for line in raw.decode("utf-8").strip().split("\n"): + line = line.strip() + if line.startswith("data:"): + line = line[5:].strip() + if not line or line == "[DONE]" or line.startswith("event:"): + continue + try: + events.append(json.loads(line)) + except json.JSONDecodeError: + continue + return events + + def _build_chunks(self): + """Build mock OpenAI-format chunks simulating Gemini tool call response.""" + # Chunk 1: text content + text_chunk = ModelResponseStream( + id="chatcmpl-123", + created=1700000000, + model="gemini-2.0-flash", + object="chat.completion.chunk", + choices=[ + StreamingChoices( + index=0, + delta=Delta(content="Let me check", role="assistant"), + finish_reason=None, + ) + ], + ) + + # Chunk 2: tool call (triggers new content block + carries args) + tool_chunk = ModelResponseStream( + id="chatcmpl-123", + created=1700000000, + model="gemini-2.0-flash", + object="chat.completion.chunk", + choices=[ + StreamingChoices( + index=0, + delta=Delta( + tool_calls=[ + ChatCompletionDeltaToolCall( + id="call_123", + type="function", + function=Function( + name="get_weather", + arguments='{"city": "Tokyo"}', + ), + index=0, + ) + ] + ), + finish_reason=None, + ) + ], + ) + + # Chunk 3: finish + finish_chunk = ModelResponseStream( + id="chatcmpl-123", + created=1700000000, + model="gemini-2.0-flash", + object="chat.completion.chunk", + choices=[ + StreamingChoices( + index=0, + delta=Delta(), + finish_reason="stop", + ) + ], + usage=Usage(prompt_tokens=10, completion_tokens=5, total_tokens=15), + ) + + return [text_chunk, tool_chunk, finish_chunk] + + def _make_stream_wrapper(self, chunks): + from litellm.llms.anthropic.experimental_pass_through.adapters.streaming_iterator import ( + AnthropicStreamWrapper, + ) + + class SimpleIterator: + def __init__(self, items): + self._items = iter(items) + + def __iter__(self): + return self + + def __next__(self): + return next(self._items) + + def __aiter__(self): + return self + + async def __anext__(self): + try: + return next(self._items) + except StopIteration: + raise StopAsyncIteration + + return AnthropicStreamWrapper( + completion_stream=SimpleIterator(chunks), + model="gemini/gemini-2.0-flash", + ) + + def _find_tool_deltas(self, events): + return [ + e for e in events + if isinstance(e, dict) + and e.get("type") == "content_block_delta" + and isinstance(e.get("delta"), dict) + and e["delta"].get("type") == "input_json_delta" + ] + + def test_sync_tool_args_not_dropped(self): + import json + + chunks = self._build_chunks() + wrapper = self._make_stream_wrapper(chunks) + + events = list(wrapper) + tool_deltas = self._find_tool_deltas(events) + + assert len(tool_deltas) > 0, ( + f"No input_json_delta events found (issue #24134). " + f"Event types: {[e.get('type') for e in events if isinstance(e, dict)]}" + ) + + combined = "".join(d["delta"]["partial_json"] for d in tool_deltas) + parsed = json.loads(combined) + assert parsed == {"city": "Tokyo"} + + @pytest.mark.asyncio + async def test_async_tool_args_not_dropped(self): + import json + + chunks = self._build_chunks() + wrapper = self._make_stream_wrapper(chunks) + + events = [] + async for event in wrapper: + events.append(event) + + tool_deltas = self._find_tool_deltas(events) + + assert len(tool_deltas) > 0, ( + f"No input_json_delta events found (issue #24134). " + f"Event types: {[e.get('type') for e in events if isinstance(e, dict)]}" + ) + + combined = "".join(d["delta"]["partial_json"] for d in tool_deltas) + parsed = json.loads(combined) + assert parsed == {"city": "Tokyo"}