fix(bugfixes): bedrock None context_mgmt; stream per-instance queue; sync polyfill; trailing-chunk passthrough

Co-authored-by: Yassin Kortam <yassin@berri.ai>
This commit is contained in:
Cursor Agent 2026-05-26 12:56:12 +00:00
parent 15ea941fbe
commit 868fa20913
No known key found for this signature in database
2 changed files with 75 additions and 66 deletions

View file

@ -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:

View file

@ -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).