diff --git a/plugins/platforms/a2a/adapter.py b/plugins/platforms/a2a/adapter.py index 73568595ff..bdb1b5fb8a 100644 --- a/plugins/platforms/a2a/adapter.py +++ b/plugins/platforms/a2a/adapter.py @@ -33,7 +33,9 @@ from . import protocol, security logger = logging.getLogger(__name__) _DEFAULT_PORT = 9900 -_MIN_ORPHAN_TIMEOUT, _WATCHDOG_INTERVAL = 300, 60 # seconds: orphan grace floor / watchdog period +# seconds: orphan grace floor / ceiling / watchdog period. The ceiling keeps the sweep +# meaningful when A2A_REPLY_TIMEOUT is absurd (1e18 would never fail an orphan). +_MIN_ORPHAN_TIMEOUT, _MAX_ORPHAN_TIMEOUT, _WATCHDOG_INTERVAL = 300, 86400, 60 _MAX_BODY = 1_048_576 # 1MB max request body — prevents DoS via memory exhaustion _SSE_KEEPALIVE = 5 # seconds between SSE keepalive comments _DEFAULT_DESCRIPTION = "Hermes Agent — a general-purpose agent reachable over A2A." @@ -70,8 +72,8 @@ def _reply_timeout() -> float: def _orphan_timeout() -> float: - """Orphan grace must never expire before a configured reply window.""" - return max(float(_MIN_ORPHAN_TIMEOUT), _reply_timeout()) + """Orphan grace must never expire before a configured reply window, but stays bounded.""" + return min(float(_MAX_ORPHAN_TIMEOUT), max(float(_MIN_ORPHAN_TIMEOUT), _reply_timeout())) def _default_agent_name() -> str: @@ -336,6 +338,7 @@ class A2AAdapter(BasePlatformAdapter): self._resolve_locked(tid, protocol.STATE_FAILED, "[agent shutting down]") self._pending.clear() self._pending_order.clear() + self._active_tasks.clear() def _watchdog_loop(self) -> None: """Background thread that fails orphaned tasks (keeps them queryable).""" diff --git a/tests/plugins/test_a2a_phase23.py b/tests/plugins/test_a2a_phase23.py index 67ac1ca66c..068b72f9ab 100644 --- a/tests/plugins/test_a2a_phase23.py +++ b/tests/plugins/test_a2a_phase23.py @@ -534,6 +534,16 @@ class TestTaskStore: adapter._pop_pending("t-live") assert adapter._fail_orphans_once() == ["t-live"] + def test_orphan_timeout_is_bounded_and_disconnect_clears_active_tasks(self, monkeypatch): + from plugins.platforms.a2a import adapter as mod + monkeypatch.setenv("A2A_REPLY_TIMEOUT", "1e18") + assert mod._orphan_timeout() == mod._MAX_ORPHAN_TIMEOUT + + adapter, _base = _make_live_adapter(monkeypatch) + adapter._add_pending("t-live", "c1") + asyncio.run(adapter.disconnect()) + assert adapter._active_tasks == set() + def test_watchdog_cannot_race_local_finalization(self, monkeypatch): adapter, _base = _make_live_adapter(monkeypatch) rec = adapter.tasks.create("t-live", "c1", "peer")