fix(agent): stream_options compatibility retry does not consume the transient budget
Review finding on #109143: `_handle_stream_error` returned True for the stream_options rejection without checking whether `_call` had another iteration. With HERMES_STREAM_RETRIES=0, or after the transient budget was spent, the loop ended with neither a response nor an error set and the call returned None instead of raising. The compatibility retry now extends the loop by one attempt exactly once (`_compat_retries`); the transient budget is untouched. Test pinned with HERMES_STREAM_RETRIES=0 (red on the previous head).
This commit is contained in:
@@ -3119,6 +3119,7 @@ class _StreamingCall(StreamingWaitMonitor):
|
||||
if not self.deltas_were_sent["yes"] and not getattr(self.agent, "_stream_options_unsupported", False) and _rejects_stream_options(e):
|
||||
# Nothing streamed yet: drop the usage extension for this session and re-open.
|
||||
self.agent._stream_options_unsupported = True
|
||||
self._compat_retries = 1
|
||||
logger.info("Endpoint rejected stream_options (HTTP %s); retrying without it for this session.",
|
||||
getattr(e, "status_code", None))
|
||||
self._cancel_current_stream_attempt("stream_options_rejected_retry")
|
||||
@@ -3176,8 +3177,14 @@ class _StreamingCall(StreamingWaitMonitor):
|
||||
|
||||
def _call(self):
|
||||
_max_stream_retries = env_int("HERMES_STREAM_RETRIES", 2)
|
||||
# The one stream_options compatibility retry (#9705) is not a network retry and must not
|
||||
# consume the transient budget: on the last attempt (or HERMES_STREAM_RETRIES=0) the
|
||||
# handler returned True and the loop ended with neither a response nor an error set.
|
||||
self._compat_retries = 0
|
||||
_stream_attempt = -1
|
||||
try:
|
||||
for _stream_attempt in range(_max_stream_retries + 1):
|
||||
while _stream_attempt < _max_stream_retries + self._compat_retries:
|
||||
_stream_attempt += 1
|
||||
stream_attempt_id = self._start_stream_attempt()
|
||||
# Otherwise /stop closes the connection and the retry opens a
|
||||
# FRESH one, blocking up to a full read timeout per attempt.
|
||||
|
||||
@@ -262,12 +262,14 @@ class TestStreamingAccumulator:
|
||||
|
||||
@patch("run_agent.AIAgent._create_request_openai_client")
|
||||
@patch("run_agent.AIAgent._close_request_openai_client")
|
||||
def test_endpoint_rejecting_stream_options_is_retried_without_it(self, mock_close, mock_create):
|
||||
def test_endpoint_rejecting_stream_options_is_retried_without_it(self, mock_close, mock_create, monkeypatch):
|
||||
"""Strict OpenAI-compatible endpoints (Azure AI Foundry MaaS) 422 on
|
||||
``stream_options.include_usage``; the call is retried once without the field and
|
||||
the session remembers the rejection (#9705)."""
|
||||
the session remembers the rejection (#9705). The compatibility retry must not spend
|
||||
the transient-retry budget: with HERMES_STREAM_RETRIES=0 it still happens."""
|
||||
from openai import APIStatusError
|
||||
from run_agent import AIAgent
|
||||
monkeypatch.setenv("HERMES_STREAM_RETRIES", "0")
|
||||
|
||||
body = {"detail": [{"type": "extra_forbidden", "loc": ["body", "stream_options", "include_usage"],
|
||||
"msg": "Extra inputs are not permitted"}]}
|
||||
|
||||
Reference in New Issue
Block a user