mirror of
https://github.com/BerriAI/litellm.git
synced 2026-10-03 02:22:24 +00:00
fix(anthropic): delay streaming content block start
This commit is contained in:
parent
c04d5e5ea9
commit
5be47aac1e
2 changed files with 57 additions and 25 deletions
|
|
@ -103,23 +103,12 @@ class AnthropicStreamWrapper(AdapterCompletionStreamWrapper):
|
|||
)
|
||||
return self.chunk_queue.popleft()
|
||||
|
||||
if self.sent_content_block_start is False:
|
||||
self.sent_content_block_start = True
|
||||
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:
|
||||
self._increment_content_block_index()
|
||||
|
||||
processed_chunk = LiteLLMAnthropicMessagesAdapter().translate_streaming_openai_response_to_anthropic(
|
||||
|
|
@ -127,6 +116,27 @@ class AnthropicStreamWrapper(AdapterCompletionStreamWrapper):
|
|||
current_content_block_index=self.current_content_block_index,
|
||||
)
|
||||
|
||||
if not self.sent_content_block_start:
|
||||
self.sent_content_block_start = True
|
||||
self.chunk_queue.append(
|
||||
{
|
||||
"type": "content_block_start",
|
||||
"index": self.current_content_block_index,
|
||||
"content_block": self.current_content_block_start,
|
||||
}
|
||||
)
|
||||
if (
|
||||
processed_chunk.get("type") == "content_block_delta"
|
||||
and isinstance(processed_chunk.get("delta"), dict)
|
||||
and processed_chunk["delta"].get("type")
|
||||
in ("text_delta", "input_json_delta")
|
||||
and processed_chunk["delta"].get(
|
||||
"text", processed_chunk["delta"].get("partial_json", "")
|
||||
)
|
||||
):
|
||||
self.chunk_queue.append(processed_chunk)
|
||||
return self.chunk_queue.popleft()
|
||||
|
||||
if should_start_new_block and not self.sent_content_block_finish:
|
||||
# Queue the sequence: content_block_stop -> content_block_start
|
||||
# For text blocks the trigger chunk is not emitted as a separate
|
||||
|
|
@ -243,24 +253,13 @@ class AnthropicStreamWrapper(AdapterCompletionStreamWrapper):
|
|||
)
|
||||
return self.chunk_queue.popleft()
|
||||
|
||||
if self.sent_content_block_start is False:
|
||||
self.sent_content_block_start = True
|
||||
self.chunk_queue.append(
|
||||
{
|
||||
"type": "content_block_start",
|
||||
"index": self.current_content_block_index,
|
||||
"content_block": {"type": "text", "text": ""},
|
||||
}
|
||||
)
|
||||
return self.chunk_queue.popleft()
|
||||
|
||||
async for chunk in self.completion_stream:
|
||||
if chunk == "None" or chunk is None:
|
||||
raise Exception
|
||||
|
||||
# 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:
|
||||
self._increment_content_block_index()
|
||||
|
||||
processed_chunk = LiteLLMAnthropicMessagesAdapter().translate_streaming_openai_response_to_anthropic(
|
||||
|
|
@ -268,6 +267,27 @@ class AnthropicStreamWrapper(AdapterCompletionStreamWrapper):
|
|||
current_content_block_index=self.current_content_block_index,
|
||||
)
|
||||
|
||||
if not self.sent_content_block_start:
|
||||
self.sent_content_block_start = True
|
||||
self.chunk_queue.append(
|
||||
{
|
||||
"type": "content_block_start",
|
||||
"index": self.current_content_block_index,
|
||||
"content_block": self.current_content_block_start,
|
||||
}
|
||||
)
|
||||
if (
|
||||
processed_chunk.get("type") == "content_block_delta"
|
||||
and isinstance(processed_chunk.get("delta"), dict)
|
||||
and processed_chunk["delta"].get("type")
|
||||
in ("text_delta", "input_json_delta")
|
||||
and processed_chunk["delta"].get(
|
||||
"text", processed_chunk["delta"].get("partial_json", "")
|
||||
)
|
||||
):
|
||||
self.chunk_queue.append(processed_chunk)
|
||||
return self.chunk_queue.popleft()
|
||||
|
||||
# Check if this is a usage chunk and we have a held stop_reason chunk
|
||||
if (
|
||||
self.holding_stop_reason_chunk is not None
|
||||
|
|
|
|||
|
|
@ -1707,13 +1707,25 @@ class AWSEventStreamDecoder:
|
|||
if "trace" in chunk_data:
|
||||
trace = chunk_data.get("trace")
|
||||
model_response_provider_specific_fields["trace"] = trace
|
||||
delta_content = text
|
||||
if tool_use is not None and delta_content == "":
|
||||
delta_content = None
|
||||
elif (
|
||||
delta_content == ""
|
||||
and tool_use is None
|
||||
and not provider_specific_fields
|
||||
and thinking_blocks is None
|
||||
and reasoning_content is None
|
||||
):
|
||||
delta_content = None
|
||||
|
||||
response = ModelResponseStream(
|
||||
choices=[
|
||||
StreamingChoices(
|
||||
finish_reason=finish_reason,
|
||||
index=0, # Always 0 - Bedrock never returns multiple choices
|
||||
delta=Delta(
|
||||
content=text,
|
||||
content=delta_content,
|
||||
role="assistant",
|
||||
tool_calls=[tool_use] if tool_use else None,
|
||||
provider_specific_fields=(
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue