mirror of
https://github.com/BerriAI/litellm.git
synced 2026-10-08 03:08:45 +00:00
Merge pull request #33315 from BerriAI/litellm_fix_empty_delta_thinking_block
fix(anthropic-adapter): drop empty content_block_delta events
This commit is contained in:
commit
03ef18a9ea
4 changed files with 178 additions and 24 deletions
|
|
@ -13,14 +13,18 @@ from typing import (
|
|||
List,
|
||||
Literal,
|
||||
Optional,
|
||||
get_args,
|
||||
)
|
||||
|
||||
from typing_extensions import assert_never
|
||||
|
||||
from litellm._logging import verbose_logger
|
||||
from litellm._uuid import uuid
|
||||
from litellm.types.llms.anthropic import (
|
||||
AppliedEdit,
|
||||
CompactionBlock,
|
||||
ContextManagementResponse,
|
||||
StreamingContentBlockDeltaType,
|
||||
UsageDelta,
|
||||
UsageIteration,
|
||||
)
|
||||
|
|
@ -30,6 +34,23 @@ if TYPE_CHECKING:
|
|||
from litellm.types.utils import ModelResponseStream
|
||||
|
||||
|
||||
_STREAMING_DELTA_TYPES = frozenset(get_args(StreamingContentBlockDeltaType))
|
||||
|
||||
|
||||
def _delta_payload_field(delta_type: StreamingContentBlockDeltaType) -> str:
|
||||
match delta_type:
|
||||
case "text_delta":
|
||||
return "text"
|
||||
case "input_json_delta":
|
||||
return "partial_json"
|
||||
case "thinking_delta":
|
||||
return "thinking"
|
||||
case "signature_delta":
|
||||
return "signature"
|
||||
case _:
|
||||
assert_never(delta_type)
|
||||
|
||||
|
||||
class _CombinedChunkSplitter:
|
||||
"""
|
||||
Splits a streaming chunk that carries BOTH response content and a
|
||||
|
|
@ -458,12 +479,15 @@ class AnthropicStreamWrapper(AdapterCompletionStreamWrapper):
|
|||
|
||||
# 3. If the trigger chunk carries delta content, queue it
|
||||
# so the first delta of the new block is not silently dropped.
|
||||
if self._trigger_delta_has_content(processed_chunk):
|
||||
if self._delta_has_content(processed_chunk):
|
||||
self.chunk_queue.append(processed_chunk)
|
||||
|
||||
self.sent_content_block_finish = False
|
||||
return self.chunk_queue.popleft()
|
||||
|
||||
if processed_chunk["type"] == "content_block_delta" and not self._delta_has_content(processed_chunk):
|
||||
continue
|
||||
|
||||
if processed_chunk["type"] == "message_delta" and self.sent_content_block_finish is False:
|
||||
# Queue both the content_block_stop and the message_delta
|
||||
self.chunk_queue.append(
|
||||
|
|
@ -670,13 +694,18 @@ class AnthropicStreamWrapper(AdapterCompletionStreamWrapper):
|
|||
|
||||
# 3. If the trigger chunk carries delta content, queue it
|
||||
# so the first delta of the new block is not silently dropped.
|
||||
if self._trigger_delta_has_content(processed_chunk):
|
||||
if self._delta_has_content(processed_chunk):
|
||||
self.chunk_queue.append(processed_chunk)
|
||||
|
||||
# Reset state for new block
|
||||
self.sent_content_block_finish = False
|
||||
return self.chunk_queue.popleft()
|
||||
|
||||
if processed_chunk["type"] == "content_block_delta" and not self._delta_has_content(
|
||||
processed_chunk
|
||||
):
|
||||
continue
|
||||
|
||||
if processed_chunk["type"] == "message_delta" and self.sent_content_block_finish is False:
|
||||
# Queue both the content_block_stop and the holding chunk
|
||||
self.chunk_queue.append(
|
||||
|
|
@ -808,20 +837,33 @@ class AnthropicStreamWrapper(AdapterCompletionStreamWrapper):
|
|||
self.current_content_block_index += 1
|
||||
|
||||
@staticmethod
|
||||
def _trigger_delta_has_content(processed_chunk: Dict[str, Any]) -> bool:
|
||||
"""Return True if a translated trigger chunk carries a non-empty
|
||||
``content_block_delta`` payload that must be re-emitted after a
|
||||
block transition.
|
||||
def _delta_has_content(processed_chunk: Dict[str, Any]) -> bool:
|
||||
"""Return True if a translated chunk carries a non-empty
|
||||
``content_block_delta`` payload.
|
||||
|
||||
When an upstream chunk both *triggers* a new content block (its type
|
||||
differs from the active block) and *carries* delta content, that
|
||||
content belongs to the new block. The synthesized
|
||||
``content_block_start`` only ever carries an empty body — see
|
||||
Gates every ``content_block_delta`` emission. An empty delta carries
|
||||
no information, and the translate fallback types empty deltas as
|
||||
``text_delta`` regardless of the active block's type — emitting one
|
||||
into an open ``thinking`` block (e.g. Bedrock Converse sends an empty
|
||||
reasoning delta mid-block) crashes strict Anthropic SDK clients with
|
||||
"Content block is not a text block".
|
||||
|
||||
Also gates re-emission after a block transition: when an upstream
|
||||
chunk both *triggers* a new content block (its type differs from the
|
||||
active block) and *carries* delta content, that content belongs to
|
||||
the new block. The synthesized ``content_block_start`` only ever
|
||||
carries an empty body — see
|
||||
``_translate_streaming_openai_chunk_to_anthropic_content_block``,
|
||||
which returns an empty ``TextBlock``/``ToolUseBlock``/thinking block —
|
||||
so the trigger chunk's delta must be re-queued or the first token of
|
||||
the new block (the first non-empty text/thinking delta, or bundled
|
||||
tool arguments) is silently dropped.
|
||||
|
||||
Delta types outside ``StreamingContentBlockDeltaType`` — the closed
|
||||
set the translate layer can produce — are treated as empty. The
|
||||
per-type payload lookup is exhaustively matched against that set in
|
||||
``_delta_payload_field``, so extending the translate layer with a new
|
||||
delta type fails type-checking here until it is handled.
|
||||
"""
|
||||
if processed_chunk.get("type") != "content_block_delta":
|
||||
return False
|
||||
|
|
@ -829,15 +871,9 @@ class AnthropicStreamWrapper(AdapterCompletionStreamWrapper):
|
|||
if not isinstance(delta, dict):
|
||||
return False
|
||||
delta_type = delta.get("type")
|
||||
if delta_type == "text_delta":
|
||||
return bool(delta.get("text"))
|
||||
if delta_type == "input_json_delta":
|
||||
return bool(delta.get("partial_json"))
|
||||
if delta_type == "thinking_delta":
|
||||
return bool(delta.get("thinking"))
|
||||
if delta_type == "signature_delta":
|
||||
return bool(delta.get("signature"))
|
||||
return False
|
||||
if delta_type not in _STREAMING_DELTA_TYPES:
|
||||
return False
|
||||
return bool(delta.get(_delta_payload_field(delta_type)))
|
||||
|
||||
def _should_start_new_content_block(self, chunk: "ModelResponseStream") -> bool:
|
||||
"""
|
||||
|
|
|
|||
|
|
@ -104,6 +104,7 @@ from litellm.types.llms.anthropic import (
|
|||
ContextManagementResponse,
|
||||
MessageBlockDelta,
|
||||
MessageDelta,
|
||||
StreamingContentBlockDeltaType,
|
||||
UsageDelta,
|
||||
UsageIteration,
|
||||
)
|
||||
|
|
@ -1423,7 +1424,7 @@ class LiteLLMAnthropicMessagesAdapter:
|
|||
def _translate_streaming_openai_chunk_to_anthropic(
|
||||
self, choices: List[Union[OpenAIStreamingChoice, StreamingChoices]]
|
||||
) -> Tuple[
|
||||
Literal["text_delta", "input_json_delta", "thinking_delta", "signature_delta"],
|
||||
StreamingContentBlockDeltaType,
|
||||
Union[
|
||||
ContentTextBlockDelta,
|
||||
ContentJsonBlockDelta,
|
||||
|
|
|
|||
|
|
@ -439,6 +439,9 @@ class ContentThinkingSignatureBlockDelta(TypedDict):
|
|||
signature: str
|
||||
|
||||
|
||||
StreamingContentBlockDeltaType = Literal["text_delta", "input_json_delta", "thinking_delta", "signature_delta"]
|
||||
|
||||
|
||||
class ContentBlockDelta(TypedDict):
|
||||
type: Literal["content_block_delta"]
|
||||
index: int
|
||||
|
|
|
|||
|
|
@ -12,6 +12,13 @@ silently dropped — e.g. text resuming after a tool call started from the secon
|
|||
token ("The weather is nice." was lost, "Hi" rendered as ""). Bundled
|
||||
``input_json_delta`` tool arguments were already preserved and must stay
|
||||
preserved, and empty trigger deltas must not produce spurious events.
|
||||
|
||||
Also covers the inverse regression: a chunk whose translated delta carries no
|
||||
payload must not be emitted at all. The translate fallback types empty deltas
|
||||
as ``text_delta`` regardless of the active block, so an empty reasoning delta
|
||||
mid-thinking-block (Bedrock Converse sends these) used to emit ``text_delta``
|
||||
into an open ``thinking`` block, crashing Anthropic SDK clients (Claude Code)
|
||||
with "Content block is not a text block".
|
||||
"""
|
||||
|
||||
import os
|
||||
|
|
@ -51,6 +58,13 @@ def _make_chunk(delta: Delta, finish_reason: Optional[str] = None) -> MagicMock:
|
|||
return chunk
|
||||
|
||||
|
||||
def _thinking_chunk(thinking: str, signature: str = "") -> MagicMock:
|
||||
block = {"type": "thinking", "thinking": thinking}
|
||||
if signature:
|
||||
block["signature"] = signature
|
||||
return _make_chunk(Delta(content=None, thinking_blocks=[block]))
|
||||
|
||||
|
||||
def _tool_chunk(
|
||||
call_id: str, name: Optional[str], arguments: Optional[str]
|
||||
) -> MagicMock:
|
||||
|
|
@ -109,6 +123,47 @@ def _input_json_deltas(events: List[dict]) -> List[str]:
|
|||
]
|
||||
|
||||
|
||||
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"
|
||||
]
|
||||
|
||||
|
||||
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"
|
||||
]
|
||||
|
||||
|
||||
_DELTA_TYPES_PER_BLOCK_TYPE = {
|
||||
"text": {"text_delta"},
|
||||
"thinking": {"thinking_delta", "signature_delta"},
|
||||
"tool_use": {"input_json_delta"},
|
||||
}
|
||||
|
||||
|
||||
def _assert_deltas_match_their_block_type(events: List[dict]) -> None:
|
||||
"""Enforce the invariant the Anthropic SDK enforces client-side: every
|
||||
``content_block_delta`` must be of a type valid for the block opened by
|
||||
the most recent ``content_block_start`` at the same index.
|
||||
"""
|
||||
block_types = {}
|
||||
for event in events:
|
||||
if event.get("type") == "content_block_start":
|
||||
block_types[event["index"]] = event["content_block"]["type"]
|
||||
if event.get("type") == "content_block_delta":
|
||||
block_type = block_types[event["index"]]
|
||||
assert event["delta"]["type"] in _DELTA_TYPES_PER_BLOCK_TYPE[block_type], (
|
||||
f"{event['delta']['type']} emitted into a {block_type} block: {event}"
|
||||
)
|
||||
|
||||
|
||||
def test_held_stop_reason_usage_merge_preserves_openai_cache_token_details():
|
||||
"""OpenAI-compatible usage chunks carry cache reads in prompt_tokens_details."""
|
||||
wrapper = AnthropicStreamWrapper(completion_stream=iter([]), model="claude-x")
|
||||
|
|
@ -332,11 +387,70 @@ def test_bundled_tool_args_on_transition_still_preserved_sync():
|
|||
({"type": "content_block_delta", "delta": None}, False),
|
||||
],
|
||||
)
|
||||
def test_trigger_delta_has_content_branches(processed_chunk, expected):
|
||||
"""Directly exercise the re-emit predicate across all delta types and the
|
||||
def test_delta_has_content_branches(processed_chunk, expected):
|
||||
"""Directly exercise the emission predicate across all delta types and the
|
||||
empty/malformed guards, so the helper's behavior is pinned independently of
|
||||
upstream chunk-translation details.
|
||||
"""
|
||||
assert (
|
||||
AnthropicStreamWrapper._trigger_delta_has_content(processed_chunk) is expected
|
||||
assert AnthropicStreamWrapper._delta_has_content(processed_chunk) is expected
|
||||
|
||||
|
||||
def _empty_reasoning_delta_mid_thinking_chunks() -> List[MagicMock]:
|
||||
return [
|
||||
_thinking_chunk("Let me think"),
|
||||
_thinking_chunk(""),
|
||||
_thinking_chunk("", signature="sig123"),
|
||||
_make_chunk(Delta(content="Hello")),
|
||||
_make_chunk(Delta(content=None), finish_reason="stop"),
|
||||
]
|
||||
|
||||
|
||||
def _assert_empty_reasoning_delta_suppressed(events: List[dict]) -> None:
|
||||
_assert_deltas_match_their_block_type(events)
|
||||
assert _thinking_deltas(events) == ["Let me think"]
|
||||
assert _signature_deltas(events) == ["sig123"]
|
||||
assert _text_deltas(events) == ["Hello"]
|
||||
|
||||
|
||||
def test_empty_reasoning_delta_mid_thinking_block_is_suppressed_sync():
|
||||
"""Bedrock Converse repro: an empty reasoning delta arriving inside an open
|
||||
thinking block used to be emitted as ``text_delta {"text": ""}`` at the
|
||||
thinking block's index (no block transition), which crashes Claude Code's
|
||||
Anthropic SDK with "Content block is not a text block". It must be dropped,
|
||||
while the surrounding thinking/signature/text deltas all still flow.
|
||||
"""
|
||||
wrapper = AnthropicStreamWrapper(
|
||||
completion_stream=iter(_empty_reasoning_delta_mid_thinking_chunks()),
|
||||
model="claude-x",
|
||||
)
|
||||
_assert_empty_reasoning_delta_suppressed(_drain_sync(wrapper))
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_empty_reasoning_delta_mid_thinking_block_is_suppressed_async():
|
||||
"""Async twin of the Bedrock Converse repro — the proxy serves the async
|
||||
iterator, so the skip must exist on that path too.
|
||||
"""
|
||||
wrapper = AnthropicStreamWrapper(
|
||||
completion_stream=_AsyncStream(_empty_reasoning_delta_mid_thinking_chunks()),
|
||||
model="claude-x",
|
||||
)
|
||||
_assert_empty_reasoning_delta_suppressed(await _drain_async(wrapper))
|
||||
|
||||
|
||||
def test_empty_content_chunk_mid_text_block_is_suppressed_sync():
|
||||
"""An empty-content chunk arriving mid-text-block (no transition) used to
|
||||
emit a pointless ``text_delta {"text": ""}``; it must be dropped without
|
||||
affecting the surrounding text deltas.
|
||||
"""
|
||||
chunks = [
|
||||
_make_chunk(Delta(content="Hi")),
|
||||
_make_chunk(Delta(content="")),
|
||||
_make_chunk(Delta(content=" there")),
|
||||
_make_chunk(Delta(content=None), finish_reason="stop"),
|
||||
]
|
||||
wrapper = AnthropicStreamWrapper(completion_stream=iter(chunks), model="claude-x")
|
||||
events = _drain_sync(wrapper)
|
||||
|
||||
assert _text_deltas(events) == ["Hi", " there"]
|
||||
_assert_deltas_match_their_block_type(events)
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue