From f15d9f807b7744907cd7706487ab650dbb761e26 Mon Sep 17 00:00:00 2001 From: milo Date: Tue, 2 Jun 2026 11:33:02 +0800 Subject: [PATCH] fix: preserve delayed Anthropic stream block start after merge --- .../adapters/streaming_iterator.py | 60 +++++++------------ 1 file changed, 22 insertions(+), 38 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 7db9c9ff351..90ff23e4c02 100644 --- a/litellm/llms/anthropic/experimental_pass_through/adapters/streaming_iterator.py +++ b/litellm/llms/anthropic/experimental_pass_through/adapters/streaming_iterator.py @@ -312,14 +312,6 @@ class AnthropicStreamWrapper(AdapterCompletionStreamWrapper): if compaction_event is not None: return compaction_event - if ( - self.sent_compaction_block is False - and self.compaction_block is not None - ): - compaction_event = self._next_compaction_event() - if compaction_event is not None: - return compaction_event - for chunk in self.completion_stream: if chunk == "None" or chunk is None: raise Exception @@ -350,28 +342,6 @@ class AnthropicStreamWrapper(AdapterCompletionStreamWrapper): ), ) - if not self.sent_content_block_start: - 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, - } - ) - 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 not self.sent_content_block_start: self.sent_content_block_start = True self.sent_content_block_finish = False @@ -599,14 +569,6 @@ class AnthropicStreamWrapper(AdapterCompletionStreamWrapper): if compaction_event is not None: return compaction_event - if ( - self.sent_compaction_block is False - and self.compaction_block is not None - ): - compaction_event = self._next_compaction_event() - if compaction_event is not None: - return compaction_event - async for chunk in self.completion_stream: if chunk == "None" or chunk is None: raise Exception @@ -638,6 +600,28 @@ class AnthropicStreamWrapper(AdapterCompletionStreamWrapper): ), ) + if not self.sent_content_block_start: + 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, + } + ) + 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 will_merge_into_held: merged_chunk = self._merge_usage_into_held_stop_reason_chunk(chunk)