Files
hermes-agent/tests/agent/test_stream_retry_backoff.py
kshitijk4poor bcf88b30c6 refactor(streaming): reuse jittered_backoff for stream reconnect delay; no SDK import in error path
Review follow-up (#60029): restart the stale clock after the backoff so the
stale monitor cannot strike a stream that has not reopened yet. the private 1/2/4 s schedule duplicated
agent.retry_utils.jittered_backoff and escaped the tests/agent autouse stub
that zeroes it, so unrelated retry tests slept for real. The Anthropic
connection-error type is now read from sys.modules (an Anthropic error implies
the SDK is loaded) instead of importing it inside the handler. Tests record the
requested delay instead of patching the global time module.
2026-09-24 22:24:48 +05:30

64 lines
2.3 KiB
Python

"""Stream-level reconnects back off exponentially and stay interruptible (#60029)."""
import threading
import time
from types import SimpleNamespace
from unittest.mock import MagicMock, patch
import httpx
import pytest
from agent import chat_completion_helpers as cch
def _bare_call(agent):
call = object.__new__(cch._StreamingCall)
call.agent = agent
call.api_kwargs = {}
call.result = {}
call.clients = SimpleNamespace(diag={}, close_once=MagicMock())
call._request_cancelled = {"value": False}
call.deltas_were_sent = {"yes": False}
call.first_delta_fired = {"done": False}
call.provider_tool_in_flight = {"yes": False}
call._cancel_current_stream_attempt = MagicMock()
call.last_chunk_time = {"t": 0.0}
return call
def _agent(**kw):
return SimpleNamespace(
_interrupt_requested=False,
_stream_options_unsupported=False,
_emit_stream_drop=MagicMock(),
_is_provider_stream_parse_error=lambda e: False,
**kw,
)
@pytest.mark.real_retry_backoff
def test_transient_drops_retry_with_capped_exponential_backoff():
"""Each reconnect waits 1s, 2s, 4s, 4s — including Anthropic SDK connection errors."""
anthropic = pytest.importorskip("anthropic")
waits = []
call = _bare_call(_agent())
req = httpx.Request("POST", "https://api.anthropic.com/v1/messages")
errors = [anthropic.APIConnectionError(request=req), httpx.ConnectError("reset"),
httpx.ReadError("reset"), httpx.RemoteProtocolError("closed")]
with patch.object(cch, "_wait_stream_retry_backoff", lambda agent, delay: waits.append(delay) or True):
for attempt, err in enumerate(errors):
assert call._handle_stream_error(err, attempt, max_retries=10) is True
assert waits == [1.0, 2.0, 4.0, 4.0]
assert call.last_chunk_time["t"] > 0.0 # stale clock restarted after the backoff
assert "error" not in call.result
@pytest.mark.real_retry_backoff
def test_backoff_wait_returns_immediately_on_interrupt():
agent = _agent()
call = _bare_call(agent)
threading.Timer(0.2, lambda: setattr(agent, "_interrupt_requested", True)).start()
start = time.monotonic()
call._retry_after_drop(httpx.ConnectError("x"), 5, 10, mid_tool_call=False, reason="t")
assert time.monotonic() - start < 1.0