mirror of
https://github.com/BerriAI/litellm.git
synced 2026-10-10 03:28:53 +00:00
fix: prevent streaming thinking blocks split into multiple content_block_start events
Addresses three distinct bugs in Anthropic pass-through adapter streaming:
1. Bug A - Eager text block creation:
- BEFORE: First content_block_start with type='text' created before examining any chunks
- AFTER: Delayed to first substantial chunk; block type determined by actual content
- IMPACT: Eliminates spurious thinking→text→thinking transitions
2. Bug B - Empty chunk transitions:
- BEFORE: No method to filter empty chunks; all classified as 'text'
- AFTER: Added _chunk_has_substantial_content() to detect meaningful chunks
- IMPACT: Empty chunks no longer trigger block transitions
3. Bug C - Unrecognized redacted_thinking:
- BEFORE: Only checks 'thinking' type, misses 'redacted_thinking'
- AFTER: Changed to in ('thinking', 'redacted_thinking')
- IMPACT: Both thinking types correctly classified as thinking blocks
All three fixes work together to maintain proper streaming state machine:
- Lazy initialization waits for first meaningful content
- Guard condition prevents empty chunks from switching blocks
- Type classification handles both thinking variants
This commit is contained in:
parent
49ca04d8c3
commit
2e9850f860
2 changed files with 60 additions and 27 deletions
|
|
@ -24,7 +24,7 @@ from litellm.types.llms.anthropic import (
|
|||
UsageDelta,
|
||||
UsageIteration,
|
||||
)
|
||||
from litellm.types.utils import AdapterCompletionStreamWrapper
|
||||
from litellm.types.utils import AdapterCompletionStreamWrapper, StreamingChoices
|
||||
|
||||
if TYPE_CHECKING:
|
||||
from litellm.types.utils import ModelResponseStream
|
||||
|
|
@ -408,26 +408,28 @@ class AnthropicStreamWrapper(AdapterCompletionStreamWrapper):
|
|||
if compaction_event is not None:
|
||||
return compaction_event
|
||||
|
||||
if self.sent_content_block_start is False:
|
||||
self.sent_content_block_start = True
|
||||
self.sent_content_block_finish = False
|
||||
self.chunk_queue.append(
|
||||
{
|
||||
"type": "content_block_start",
|
||||
"index": self.current_content_block_index,
|
||||
"content_block": {"type": "text", "text": ""},
|
||||
}
|
||||
)
|
||||
return self.chunk_queue.popleft()
|
||||
|
||||
for chunk in self.completion_stream:
|
||||
if chunk == "None" or chunk is None:
|
||||
raise Exception
|
||||
|
||||
should_start_new_block = self._should_start_new_content_block(chunk)
|
||||
if should_start_new_block:
|
||||
if should_start_new_block and self.sent_content_block_start is True:
|
||||
self._increment_content_block_index()
|
||||
|
||||
if self.sent_content_block_start is False:
|
||||
if not self._chunk_has_substantial_content(chunk):
|
||||
continue
|
||||
self.sent_content_block_start = True
|
||||
self.sent_content_block_finish = False
|
||||
self.chunk_queue.append(
|
||||
{
|
||||
"type": "content_block_start",
|
||||
"index": self.current_content_block_index,
|
||||
"content_block": self.current_content_block_start,
|
||||
}
|
||||
)
|
||||
return self.chunk_queue.popleft()
|
||||
|
||||
# applied_edits only needs to flow to the final message_delta
|
||||
# (when finish_reason is set); skip threading it through every
|
||||
# intermediate chunk. For the hold-and-merge path below,
|
||||
|
|
@ -656,16 +658,7 @@ class AnthropicStreamWrapper(AdapterCompletionStreamWrapper):
|
|||
return compaction_event
|
||||
|
||||
if self.sent_content_block_start is False:
|
||||
self.sent_content_block_start = True
|
||||
self.sent_content_block_finish = False
|
||||
self.chunk_queue.append(
|
||||
{
|
||||
"type": "content_block_start",
|
||||
"index": self.current_content_block_index,
|
||||
"content_block": {"type": "text", "text": ""},
|
||||
}
|
||||
)
|
||||
return self.chunk_queue.popleft()
|
||||
pass
|
||||
|
||||
async for chunk in self.completion_stream:
|
||||
if chunk == "None" or chunk is None:
|
||||
|
|
@ -673,9 +666,23 @@ class AnthropicStreamWrapper(AdapterCompletionStreamWrapper):
|
|||
|
||||
# Check if we need to start a new content block
|
||||
should_start_new_block = self._should_start_new_content_block(chunk)
|
||||
if should_start_new_block:
|
||||
if should_start_new_block and self.sent_content_block_start is True:
|
||||
self._increment_content_block_index()
|
||||
|
||||
if self.sent_content_block_start is False:
|
||||
if not self._chunk_has_substantial_content(chunk):
|
||||
continue
|
||||
self.sent_content_block_start = True
|
||||
self.sent_content_block_finish = False
|
||||
self.chunk_queue.append(
|
||||
{
|
||||
"type": "content_block_start",
|
||||
"index": self.current_content_block_index,
|
||||
"content_block": self.current_content_block_start,
|
||||
}
|
||||
)
|
||||
return self.chunk_queue.popleft()
|
||||
|
||||
# applied_edits only needs to flow to the final message_delta
|
||||
# (when finish_reason is set); skip threading it through every
|
||||
# intermediate chunk. For the hold-and-merge path below,
|
||||
|
|
@ -939,6 +946,12 @@ class AnthropicStreamWrapper(AdapterCompletionStreamWrapper):
|
|||
tool_block["name"] = original_name
|
||||
|
||||
if block_type != self.current_content_block_type:
|
||||
if (
|
||||
block_type == "text"
|
||||
and self.current_content_block_type == "thinking"
|
||||
and not self._chunk_has_substantial_content(chunk)
|
||||
):
|
||||
return False
|
||||
self.current_content_block_type = block_type
|
||||
self.current_content_block_start = content_block_start
|
||||
return True
|
||||
|
|
@ -957,3 +970,23 @@ class AnthropicStreamWrapper(AdapterCompletionStreamWrapper):
|
|||
return True
|
||||
|
||||
return False
|
||||
|
||||
@staticmethod
|
||||
def _chunk_has_substantial_content(chunk: "ModelResponseStream") -> bool:
|
||||
"""Return True when the chunk carries content that should determine
|
||||
the initial content block type (non-empty text, tool_calls,
|
||||
thinking_blocks, or reasoning_content). Role-only chunks and empty
|
||||
deltas return False."""
|
||||
for choice in chunk.choices:
|
||||
if not hasattr(choice, "delta"):
|
||||
continue
|
||||
if choice.delta.tool_calls is not None and len(choice.delta.tool_calls) > 0:
|
||||
return True
|
||||
if getattr(choice.delta, "content", None) and len(choice.delta.content) > 0:
|
||||
return True
|
||||
if isinstance(choice, StreamingChoices):
|
||||
if hasattr(choice.delta, "thinking_blocks") and choice.delta.thinking_blocks:
|
||||
return True
|
||||
if getattr(choice.delta, "reasoning_content", None):
|
||||
return True
|
||||
return False
|
||||
|
|
|
|||
|
|
@ -1495,7 +1495,7 @@ class LiteLLMAnthropicMessagesAdapter:
|
|||
thinking_blocks = choice.delta.thinking_blocks or []
|
||||
if len(thinking_blocks) > 0:
|
||||
thinking_block = thinking_blocks[0]
|
||||
if thinking_block["type"] == "thinking":
|
||||
if thinking_block["type"] in ("thinking", "redacted_thinking"):
|
||||
thinking = thinking_block.get("thinking") or ""
|
||||
signature = thinking_block.get("signature") or ""
|
||||
|
||||
|
|
@ -1556,7 +1556,7 @@ class LiteLLMAnthropicMessagesAdapter:
|
|||
thinking_blocks = choice.delta.thinking_blocks or []
|
||||
if len(thinking_blocks) > 0:
|
||||
for thinking_block in thinking_blocks:
|
||||
if thinking_block["type"] == "thinking":
|
||||
if thinking_block["type"] in ("thinking", "redacted_thinking"):
|
||||
thinking = thinking_block.get("thinking") or ""
|
||||
signature = thinking_block.get("signature") or ""
|
||||
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue