diff --git a/plugins/platforms/a2a/adapter.py b/plugins/platforms/a2a/adapter.py index 79842c88c676..e1fafa815730 100644 --- a/plugins/platforms/a2a/adapter.py +++ b/plugins/platforms/a2a/adapter.py @@ -60,7 +60,7 @@ logger = logging.getLogger(__name__) _DEFAULT_PORT = 9900 -_ORPHAN_TIMEOUT = 300 # seconds before a pending task is considered orphaned +_MIN_ORPHAN_TIMEOUT = 300 # preserve the existing five-minute floor _WATCHDOG_INTERVAL = 60 # seconds between orphaned task watchdog runs _MAX_BODY = 1_048_576 # 1MB max request body — prevents DoS via memory exhaustion _SSE_KEEPALIVE = 5 # seconds between SSE keepalive comments @@ -74,6 +74,16 @@ def _reply_timeout() -> float: return 300.0 +def _orphan_timeout() -> float: + """Keep pending tasks alive through the configured reply deadline. + + The watchdog runs independently from request waiters, so add one watchdog + interval of grace. Otherwise a reply timeout above five minutes is + ineffective because the task is marked failed at the old fixed deadline. + """ + return max(float(_MIN_ORPHAN_TIMEOUT), _reply_timeout() + _WATCHDOG_INTERVAL) + + def _default_agent_name() -> str: name = os.getenv("A2A_AGENT_NAME", "").strip() if name: @@ -473,8 +483,9 @@ def _watchdog_loop(self) -> None: """Background thread that fails orphaned tasks (keeps them queryable).""" while not self._watchdog_stop.wait(_WATCHDOG_INTERVAL): try: - for tid in self.tasks.fail_orphans(_ORPHAN_TIMEOUT): - logger.warning("A2A: orphaned task %s marked failed (timeout %ds)", tid, _ORPHAN_TIMEOUT) + orphan_timeout = _orphan_timeout() + for tid in self.tasks.fail_orphans(orphan_timeout): + logger.warning("A2A: orphaned task %s marked failed (timeout %.0fs)", tid, orphan_timeout) protocol.metrics.tasks_failed += 1 except Exception: logger.debug("A2A: watchdog error", exc_info=True) diff --git a/tests/plugins/test_a2a_plugin.py b/tests/plugins/test_a2a_plugin.py index 346284b27796..e66d2e54df5f 100644 --- a/tests/plugins/test_a2a_plugin.py +++ b/tests/plugins/test_a2a_plugin.py @@ -35,6 +35,26 @@ def _free_port() -> int: return port +# -------------------------------------------------------------------------- +# Timeouts +# -------------------------------------------------------------------------- + +class TestA2ATimeouts: + def test_orphan_timeout_never_preempts_configured_reply_timeout(self, monkeypatch): + from plugins.platforms.a2a import adapter + + monkeypatch.setenv("A2A_REPLY_TIMEOUT", "900") + assert adapter._orphan_timeout() >= ( + adapter._reply_timeout() + adapter._WATCHDOG_INTERVAL + ) + + def test_orphan_timeout_keeps_existing_five_minute_floor(self, monkeypatch): + from plugins.platforms.a2a import adapter + + monkeypatch.setenv("A2A_REPLY_TIMEOUT", "1") + assert adapter._orphan_timeout() >= adapter._MIN_ORPHAN_TIMEOUT + + # -------------------------------------------------------------------------- # Security # --------------------------------------------------------------------------