fix(streaming): switch main turn to non-streaming on out-of-order Anthropic SSE (#72833)

_maybe_disable_streaming now reuses anthropic_adapter._is_stream_unavailable_error
(its stream-unsupported + Bedrock IAM checks were a verbatim duplicate), so a custom
anthropic_messages provider raising 'Unexpected event order' retries without
streaming instead of repeating the same broken stream. Bedrock keeps
turn_recovery's sticky Converse switch.

Also drive the #60683 usage:null test through _call_anthropic so dropping the
normalize_stream_usage wiring in _open_anthropic_stream turns it red.
This commit is contained in:
kshitijk4poor
2026-09-24 18:19:37 +05:30
committed by kshitij
parent 2cd25c9094
commit 44f6d3bbee
2 changed files with 47 additions and 31 deletions

View File

@@ -3449,7 +3449,8 @@ class _StreamingCall(StreamingWaitMonitor):
def _maybe_disable_streaming(self, e) -> None:
"""Flip to non-streaming for failures streaming itself cannot survive, or that
re-streaming can only repeat: the provider rejecting streams outright,
AnthropicBedrock IAM lacking InvokeModelWithResponseStream, or a gateway answering
AnthropicBedrock IAM lacking InvokeModelWithResponseStream, a custom anthropic_messages
provider emitting SSE events out of order (#72833), or a gateway answering
with contentless SSE keepalive frames (a degraded gateway answers every
streaming request that way, so the retry must change channel to make progress)."""
if _is_provider_stream_empty_frame_error(e):
@@ -3465,23 +3466,32 @@ class _StreamingCall(StreamingWaitMonitor):
"⚠️ Provider stream returned an empty keepalive frame — retrying this turn "
"without streaming (streaming stays off for this session).")
return
from agent.anthropic_adapter import _is_stream_unavailable_error
if not _is_stream_unavailable_error(e):
return
_err_lower = str(e).lower()
_is_stream_unsupported = "stream" in _err_lower and "not supported" in _err_lower
_is_bedrock_stream_denied = False
if not _is_stream_unsupported and "invokemodelwithresponsestream" in _err_lower:
# Message pre-check first: importing bedrock_adapter triggers a lazy boto3 install.
from agent.bedrock_adapter import is_streaming_access_denied_error
_is_bedrock_stream_denied = is_streaming_access_denied_error(e)
if _is_stream_unsupported or _is_bedrock_stream_denied:
if "unexpected event order" in _err_lower and not _is_stream_unsupported:
# Custom anthropic_messages SSE out of order (#72833): re-streaming repeats it.
# Bedrock keeps turn_recovery's sticky Converse switch instead.
if self.agent.api_mode != "anthropic_messages" or self.agent.provider == "bedrock":
return
self.agent._disable_streaming = True
self.agent._safe_print(
"\n⚠ AWS IAM denied bedrock:InvokeModelWithResponseStream. Switching to non-streaming.\n"
" Grant that action to restore streaming output.\n"
if _is_bedrock_stream_denied else
"\n⚠ Streaming is not supported for this model/provider. Switching to non-streaming.\n"
" To avoid this delay, set display.streaming: false in config.yaml\n",
"\n⚠ Provider sent Anthropic stream events out of order. Switching to non-streaming.\n",
diagnostic=True,
)
return
# Remaining matches: stream rejected outright, or Bedrock IAM stream denial.
self.agent._disable_streaming = True
self.agent._safe_print(
"\n⚠ Streaming is not supported for this model/provider. Switching to non-streaming.\n"
" To avoid this delay, set display.streaming: false in config.yaml\n"
if _is_stream_unsupported else
"\n⚠ AWS IAM denied bedrock:InvokeModelWithResponseStream. Switching to non-streaming.\n"
" Grant that action to restore streaming output.\n",
diagnostic=True,
)
def _handle_stream_error(self, e: Exception, attempt: int, max_retries: int) -> bool:
"""Classify a failed attempt: True = retry; False = stop with

View File

@@ -95,26 +95,32 @@ def test_null_usage_aux_stream_is_normalized_without_create_retry():
assert created == []
def _anthropic_agent(stream_factory, provider="custom"):
from run_agent import AIAgent
agent = AIAgent(api_key="k", base_url="https://api.minimax.io/anthropic", model="MiniMax-M2",
quiet_mode=True, skip_context_files=True, skip_memory=True)
agent.api_mode, agent.provider, agent._interrupt_requested = "anthropic_messages", provider, False
client = SimpleNamespace(messages=SimpleNamespace(stream=lambda **kw: stream_factory()))
agent._anthropic_client = client
agent._create_request_anthropic_client = lambda *a, **k: client
return agent
def test_null_usage_stream_is_normalized_before_sdk_accumulation():
# #60683 main turn: MiniMax sends usage:null on message_start/message_delta; the SDK's
# accumulate_event() crashes mid-iteration unless _call_anthropic normalizes the raw events.
from anthropic import NOT_GIVEN
from anthropic._models import construct_type_unchecked
from anthropic.lib.streaming import MessageStream
from anthropic.types import RawMessageStreamEvent
from agent.anthropic_adapter import normalize_stream_usage
raw = [construct_type_unchecked(type_=RawMessageStreamEvent, value=v) for v in (
{"type": "message_start", "message": {"id": "m", "type": "message", "role": "assistant",
"model": "MiniMax-M2", "content": [], "usage": None}},
{"type": "content_block_start", "index": 0, "content_block": {"type": "text", "text": ""}},
{"type": "content_block_delta", "index": 0, "delta": {"type": "text_delta", "text": "hi"}},
{"type": "content_block_stop", "index": 0},
{"type": "message_delta", "delta": {"stop_reason": "end_turn"}, "usage": None},
{"type": "message_stop"},
)]
stream = normalize_stream_usage(MessageStream(iter(raw), output_format=NOT_GIVEN))
list(stream)
final = stream.get_final_message()
agent = _anthropic_agent(lambda: _Ctx(_minimax_sdk_stream()))
final = agent._interruptible_streaming_api_call({"model": "MiniMax-M2", "messages": [], "max_tokens": 8})
assert final.content[0].text == "hi" and final.stop_reason == "end_turn"
@pytest.mark.parametrize("provider,disabled", [("custom", True), ("bedrock", False)])
def test_main_turn_event_order_error_disables_streaming(provider, disabled):
# #72833 main turn: a custom anthropic_messages provider's out-of-order SSE must switch the
# retry to non-streaming; Bedrock keeps turn_recovery's Converse fallback instead.
err = RuntimeError('Unexpected event order, got content_block_delta before "message_start"')
agent = _anthropic_agent(lambda: _BrokenStream(err), provider=provider)
with pytest.raises(RuntimeError, match="Unexpected event order"):
agent._interruptible_streaming_api_call({"model": "MiniMax-M2", "messages": [], "max_tokens": 8})
assert bool(getattr(agent, "_disable_streaming", False)) is disabled