fix(moa): stream completed aggregator responses safely
Convert a completed MoA aggregator response into one valid Chat Completions delta chunk at the MoA facade boundary, normalize completed message.tool_calls into indexed stream deltas, and classify these local MoA adapter-shape errors as non-fallback format errors so a local compatibility bug cannot silently drift the user's MoA route to a single model (#55933 follow-up). Salvaged from PR #74903 by @liusencomic-cyber.
This commit is contained in:
@@ -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,
|
||||
|
||||
@@ -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,
|
||||
)
|
||||
|
||||
@@ -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.
|
||||
|
||||
@@ -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:
|
||||
|
||||
@@ -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."""
|
||||
|
||||
Reference in New Issue
Block a user