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
This commit is contained in:
Chesars 2026-03-22 10:33:21 -03:00
parent f5194b5ce3
commit 84b8c652e4
2 changed files with 187 additions and 12 deletions

View file

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

View file

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