From 5be47aac1e70cac1ec55e21b167cd8b6245133fe Mon Sep 17 00:00:00 2001 From: milo Date: Fri, 22 May 2026 22:27:26 +0800 Subject: [PATCH] fix(anthropic): delay streaming content block start --- .../adapters/streaming_iterator.py | 68 ++++++++++++------- litellm/llms/bedrock/chat/invoke_handler.py | 14 +++- 2 files changed, 57 insertions(+), 25 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 c65dfb22730..79c8a39ec5c 100644 --- a/litellm/llms/anthropic/experimental_pass_through/adapters/streaming_iterator.py +++ b/litellm/llms/anthropic/experimental_pass_through/adapters/streaming_iterator.py @@ -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 diff --git a/litellm/llms/bedrock/chat/invoke_handler.py b/litellm/llms/bedrock/chat/invoke_handler.py index 7a9916f1f31..6c6f5ef0ce0 100644 --- a/litellm/llms/bedrock/chat/invoke_handler.py +++ b/litellm/llms/bedrock/chat/invoke_handler.py @@ -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=(