diff --git a/agent/auxiliary_client.py b/agent/auxiliary_client.py index 4f8ee48de5..c565628152 100644 --- a/agent/auxiliary_client.py +++ b/agent/auxiliary_client.py @@ -8110,6 +8110,14 @@ def call_llm( kwargs["stream"] = True if stream_options: kwargs["stream_options"] = stream_options + if task == "moa_aggregator": + # MoA's own facade owns the streaming contract. Routing this nested + # aggregator stream through Relay's generic stream manager can make + # Codex Responses adapters return a completed SimpleNamespace that + # the manager then tries to iterate (#55933). Return the provider + # call directly; the MoA facade converts a completed response into + # a one-chunk delta iterator at its boundary. + return client.chat.completions.create(**kwargs) return _relay_sync_stream( client, kwargs, diff --git a/agent/chat_completion_helpers.py b/agent/chat_completion_helpers.py index 1c49a59e2a..d11b830fd8 100644 --- a/agent/chat_completion_helpers.py +++ b/agent/chat_completion_helpers.py @@ -3979,7 +3979,7 @@ def interruptible_streaming_api_call(agent, api_kwargs: dict, *, on_first_delta= " To avoid this delay, set display.streaming: false " "in config.yaml\n" ) - logger.info( + logger.exception( "Streaming failed before delivery: %s", e, ) diff --git a/agent/error_classifier.py b/agent/error_classifier.py index 27df08f58e..8ac0b6c872 100644 --- a/agent/error_classifier.py +++ b/agent/error_classifier.py @@ -836,6 +836,19 @@ def classify_api_error( if classified is not None: return classified + # Local MoA streaming compatibility errors are adapter-shape bugs, not a + # provider outage. Falling back to another model would silently switch the + # user's selected MoA route to a single-model answer (#55933 follow-up). + if provider_lower == "moa" and ( + "'types.SimpleNamespace' object is not iterable" in str(error) + or "'types.SimpleNamespace' object has no attribute 'index'" in str(error) + ): + return _result( + FailoverReason.format_error, + retryable=False, + should_fallback=False, + ) + # Local MoA config drift is deterministic: a persisted session can retain # a preset name that was later renamed/deleted. Retrying the same lookup # cannot recover and makes a clear config error look like an API outage. diff --git a/agent/moa_loop.py b/agent/moa_loop.py index 173816c8cb..f302fe3813 100644 --- a/agent/moa_loop.py +++ b/agent/moa_loop.py @@ -13,6 +13,7 @@ import logging import re import threading from concurrent.futures import ThreadPoolExecutor, wait as _futures_wait +from types import SimpleNamespace from typing import Any from agent.auxiliary_client import call_llm @@ -1300,6 +1301,53 @@ def aggregate_moa_context( ) +def _completed_response_as_stream_chunk(response: Any) -> Any: + """Convert a completed Chat Completions response into one delta stream chunk. + + MoA's outer streaming consumer expects ``choices[0].delta`` chunks. A + completed aggregator response carries ``choices[0].message`` instead; adapt + it here, at the MoA facade boundary, so provider-specific Relay behavior and + other transports remain untouched. + """ + + choices = getattr(response, "choices", None) + first_choice = choices[0] if isinstance(choices, (list, tuple)) and choices else None + message = getattr(first_choice, "message", None) + raw_tool_calls = getattr(message, "tool_calls", None) + tool_call_deltas = None + if isinstance(raw_tool_calls, (list, tuple)) and raw_tool_calls: + tool_call_deltas = [] + for index, tc in enumerate(raw_tool_calls): + function = getattr(tc, "function", None) + tool_call_deltas.append(SimpleNamespace( + index=getattr(tc, "index", index), + id=getattr(tc, "id", None), + type=getattr(tc, "type", None) or "function", + function=SimpleNamespace( + name=getattr(function, "name", None), + arguments=getattr(function, "arguments", None), + ), + )) + delta = SimpleNamespace( + content=getattr(message, "content", None), + tool_calls=tool_call_deltas, + reasoning_content=getattr(message, "reasoning_content", None), + reasoning=getattr(message, "reasoning", None), + reasoning_details=getattr(message, "reasoning_details", None), + ) + choice = SimpleNamespace( + index=getattr(first_choice, "index", 0), + delta=delta, + finish_reason=getattr(first_choice, "finish_reason", None) or "stop", + ) + return SimpleNamespace( + id=getattr(response, "id", None), + model=getattr(response, "model", None), + choices=[choice], + usage=getattr(response, "usage", None), + ) + + def _attach_reference_guidance(agg_messages: list[dict[str, Any]], guidance: str) -> None: """Attach the per-turn reference block at the END of the aggregator prompt. @@ -1685,6 +1733,14 @@ class MoAChatCompletions: self._pending_trace["aggregator_output"] = _extract_text(_agg_response) except Exception: # pragma: no cover - defensive self._pending_trace["aggregator_output"] = None + if stream and hasattr(_agg_response, "choices"): + # Some aggregator adapters (notably openai-codex Responses) consume + # their provider stream internally and return a completed response + # object even when the acting consumer requested token streaming. + # The outer chat-completions streaming loop expects delta chunks; + # hand it a one-chunk iterator instead of letting it iterate the + # SimpleNamespace response itself (#55933). + return iter((_completed_response_as_stream_chunk(_agg_response),)) return _agg_response def create(self, **api_kwargs: Any) -> Any: diff --git a/tests/run_agent/test_moa_streaming.py b/tests/run_agent/test_moa_streaming.py index dff87c531e..f40822e2b3 100644 --- a/tests/run_agent/test_moa_streaming.py +++ b/tests/run_agent/test_moa_streaming.py @@ -90,6 +90,46 @@ def test_create_streams_aggregator_when_requested(monkeypatch, tmp_path): assert agg["tools"] is not None +def test_create_wraps_completed_aggregator_response_as_delta_chunk(monkeypatch, tmp_path): + """When an aggregator adapter returns a completed response despite + stream=True (Codex Responses compatibility shape), MoA must return a + one-chunk delta iterator for the outer streaming accumulator instead of + the raw non-iterable response object (#55933). + """ + + completed = _response("aggregator acted") + completed.choices[0].message.tool_calls = [ + SimpleNamespace( + id="call_1", + type="function", + function=SimpleNamespace(name="read_file", arguments='{"path":"x"}'), + ) + ] + + def on_call(kwargs): + if kwargs["task"] == "moa_aggregator": + return completed + return None + + facade, calls = _facade(monkeypatch, tmp_path, on_call=on_call) + stream = facade.create( + messages=[{"role": "user", "content": "q"}], + tools=[], + stream=True, + ) + + chunk = next(iter(stream)) + assert chunk.choices[0].delta.content == "aggregator acted" + assert chunk.choices[0].delta.tool_calls[0].index == 0 + assert chunk.choices[0].delta.tool_calls[0].function.name == "read_file" + assert chunk.choices[0].finish_reason == "stop" + with pytest.raises(StopIteration): + next(stream) + + agg = next(c for c in calls if c["task"] == "moa_aggregator") + assert agg["stream"] is True + + def test_create_non_stream_path_unchanged(monkeypatch, tmp_path): """Default (no stream): the aggregator call carries NO stream/stream_options keys, so the non-streaming path is byte-identical to before."""