fix: bound A2A orphan grace and clear _active_tasks on disconnect
`_orphan_timeout()` was `max(300, A2A_REPLY_TIMEOUT)` with no ceiling, so an absurd value (1e18) meant the watchdog sweep could never fail an orphan — the reply window is a floor for the grace, not a licence to disable the sweep. Cap it at 86400s. `disconnect()` failed and cleared `_pending`/`_pending_order` but left `_active_tasks` populated, so a reconnected adapter would keep excluding dead task ids from the orphan sweep forever. Clear it in the same locked block.
This commit is contained in:
@@ -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)."""
|
||||
|
||||
@@ -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")
|
||||
|
||||
Reference in New Issue
Block a user