From 868fa20913fe5d3cc9ed2026edc240dec33de761 Mon Sep 17 00:00:00 2001 From: Cursor Agent Date: Tue, 26 May 2026 12:56:12 +0000 Subject: [PATCH] fix(bugfixes): bedrock None context_mgmt; stream per-instance queue; sync polyfill; trailing-chunk passthrough Co-authored-by: Yassin Kortam --- .../adapters/streaming_iterator.py | 129 +++++++++--------- .../bedrock/chat/converse_transformation.py | 12 +- 2 files changed, 75 insertions(+), 66 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 3fdd2e283ca..b51a3c69d65 100644 --- a/litellm/llms/anthropic/experimental_pass_through/adapters/streaming_iterator.py +++ b/litellm/llms/anthropic/experimental_pass_through/adapters/streaming_iterator.py @@ -54,7 +54,6 @@ class AnthropicStreamWrapper(AdapterCompletionStreamWrapper): type="text", text="", ) - chunk_queue: deque = deque() # Queue for buffering multiple chunks def __init__( self, @@ -69,6 +68,10 @@ class AnthropicStreamWrapper(AdapterCompletionStreamWrapper): self.tool_name_mapping = tool_name_mapping or {} # Polyfill applied_edits on final message_delta. self.applied_edits: List[AppliedEdit] = list(applied_edits or []) + # Per-instance queue for buffering multiple chunks. Must be initialized + # here (not at class level) so concurrent streams don't share the same + # deque and corrupt each other's SSE event order. + self.chunk_queue: deque = deque() def _create_initial_usage_delta(self) -> UsageDelta: """ @@ -91,7 +94,7 @@ class AnthropicStreamWrapper(AdapterCompletionStreamWrapper): cache_read_input_tokens=0, ) - def __next__(self): + def __next__(self): # noqa: PLR0915 from .transformation import LiteLLMAnthropicMessagesAdapter try: @@ -204,75 +207,77 @@ class AnthropicStreamWrapper(AdapterCompletionStreamWrapper): self.holding_stop_reason_chunk = None return self.chunk_queue.popleft() - if not self.queued_usage_chunk: - 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 - # delta because content_block_start carries the information. - # For tool_use blocks we must also emit the trigger chunk's delta - # when it carries input_json_delta data, because some providers - # (e.g. xAI, Gemini) include tool arguments in the same streaming - # chunk as the function name/id. + if self.queued_usage_chunk: + # Usage has already been merged + emitted. Pass any trailing + # provider events through directly instead of silently + # dropping them in the for-loop body. + self.chunk_queue.append(processed_chunk) + return self.chunk_queue.popleft() - # 1. Stop current content block - self.chunk_queue.append( - { - "type": "content_block_stop", - "index": max(self.current_content_block_index - 1, 0), - } - ) + 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 + # delta because content_block_start carries the information. + # For tool_use blocks we must also emit the trigger chunk's delta + # when it carries input_json_delta data, because some providers + # (e.g. xAI, Gemini) include tool arguments in the same streaming + # chunk as the function name/id. - # 2. Start new content block - self.chunk_queue.append( - { - "type": "content_block_start", - "index": self.current_content_block_index, - "content_block": self.current_content_block_start, - } - ) + # 1. Stop current content block + self.chunk_queue.append( + { + "type": "content_block_stop", + "index": max(self.current_content_block_index - 1, 0), + } + ) - # 3. If the trigger chunk carries tool argument data, queue it - # so the input_json_delta is not silently dropped. - if ( - processed_chunk.get("type") == "content_block_delta" - and isinstance(processed_chunk.get("delta"), dict) - and processed_chunk["delta"].get("type") - == "input_json_delta" - and processed_chunk["delta"].get("partial_json") - ): - self.chunk_queue.append(processed_chunk) - - self.sent_content_block_finish = False - return self.chunk_queue.popleft() + # 2. Start new content block + self.chunk_queue.append( + { + "type": "content_block_start", + "index": self.current_content_block_index, + "content_block": self.current_content_block_start, + } + ) + # 3. If the trigger chunk carries tool argument data, queue it + # so the input_json_delta is not silently dropped. if ( - processed_chunk["type"] == "message_delta" - and self.sent_content_block_finish is False + processed_chunk.get("type") == "content_block_delta" + and isinstance(processed_chunk.get("delta"), dict) + and processed_chunk["delta"].get("type") == "input_json_delta" + and processed_chunk["delta"].get("partial_json") ): - # Queue both the content_block_stop and the message_delta - self.chunk_queue.append( - { - "type": "content_block_stop", - "index": self.current_content_block_index, - } - ) - self.sent_content_block_finish = True - if ( - processed_chunk.get("delta", {}).get("stop_reason") - is not None - ): - self.holding_stop_reason_chunk = processed_chunk - else: - self.chunk_queue.append(processed_chunk) - return self.chunk_queue.popleft() - elif self.holding_chunk is not None: - self.chunk_queue.append(self.holding_chunk) self.chunk_queue.append(processed_chunk) - self.holding_chunk = None - return self.chunk_queue.popleft() + + self.sent_content_block_finish = False + return self.chunk_queue.popleft() + + 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( + { + "type": "content_block_stop", + "index": self.current_content_block_index, + } + ) + self.sent_content_block_finish = True + if processed_chunk.get("delta", {}).get("stop_reason") is not None: + self.holding_stop_reason_chunk = processed_chunk else: self.chunk_queue.append(processed_chunk) - return self.chunk_queue.popleft() + return self.chunk_queue.popleft() + elif self.holding_chunk is not None: + self.chunk_queue.append(self.holding_chunk) + self.chunk_queue.append(processed_chunk) + self.holding_chunk = None + return self.chunk_queue.popleft() + else: + self.chunk_queue.append(processed_chunk) + return self.chunk_queue.popleft() # Handle any remaining held chunks after stream ends if not self.queued_usage_chunk: diff --git a/litellm/llms/bedrock/chat/converse_transformation.py b/litellm/llms/bedrock/chat/converse_transformation.py index eaeb017b7c2..370d75780d7 100644 --- a/litellm/llms/bedrock/chat/converse_transformation.py +++ b/litellm/llms/bedrock/chat/converse_transformation.py @@ -979,11 +979,15 @@ class AmazonConverseConfig(BaseConfig): def _map_context_management_param( self, value: Union[dict, list], optional_params: dict ) -> None: - optional_params["context_management"] = ( - AnthropicConfig.map_openai_context_management_to_anthropic( - cast(Union[dict, list], value) - ) + mapped = AnthropicConfig.map_openai_context_management_to_anthropic( + cast(Union[dict, list], value) ) + # Skip when the mapper returned None for malformed input — leaving the + # key out is safer than passing `context_management: null` downstream, + # which Bedrock would reject and which can confuse intermediate checks + # before the final _filter_context_management_for_bedrock_converse step. + if mapped is not None: + optional_params["context_management"] = mapped def _map_service_tier_param(self, value: str, optional_params: dict) -> None: """Map OpenAI service_tier (string) to Bedrock serviceTier (object).