From 44f6d3bbeee093ab985bb1e16d7346fdc3b1db7a Mon Sep 17 00:00:00 2001 From: kshitijk4poor <82637225+kshitijk4poor@users.noreply.github.com> Date: Thu, 24 Sep 2026 18:19:37 +0530 Subject: [PATCH] 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. --- agent/chat_completion_helpers.py | 34 +++++++++----- .../agent/test_anthropic_stream_fallbacks.py | 44 +++++++++++-------- 2 files changed, 47 insertions(+), 31 deletions(-) diff --git a/agent/chat_completion_helpers.py b/agent/chat_completion_helpers.py index 18b41580d9..5b6307042f 100644 --- a/agent/chat_completion_helpers.py +++ b/agent/chat_completion_helpers.py @@ -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 diff --git a/tests/agent/test_anthropic_stream_fallbacks.py b/tests/agent/test_anthropic_stream_fallbacks.py index 86097fed71..2c57b5c4ad 100644 --- a/tests/agent/test_anthropic_stream_fallbacks.py +++ b/tests/agent/test_anthropic_stream_fallbacks.py @@ -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