Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion plugins/platforms/a2a/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -83,7 +83,7 @@ via `tasks/get`.
| `A2A_ALLOW_ALL_USERS` | `false` | Allow any authed peer (dev only). |
| `A2A_RATE_LIMIT` | `60` | Requests/minute per identity. |
| `A2A_MAX_PINGPONG_TURNS` | `5` | Anti-loop turn cap per context (max 20). |
| `A2A_REPLY_TIMEOUT` | `300` | Seconds to wait for the agent's reply; the orphan sweep never fails a task before this window (floor 300s) or while a request still waits on it. |
| `A2A_REPLY_TIMEOUT` | `300` | Seconds a request waits for the agent's reply; past it (or when a stream client disconnects) the task stays working and the reply is stored when it arrives, up to the orphan ceiling (24 h). |
| `A2A_PUSH_SECRET` | bearer token | HMAC secret for push signing. |
| `A2A_ADVERTISED_TOOLSETS` | all registered | Restrict skills on the Agent Card. |

Expand Down
49 changes: 41 additions & 8 deletions plugins/platforms/a2a/adapter.py
Original file line number Diff line number Diff line change
Expand Up @@ -10,6 +10,7 @@
import json
import logging
import os
import socket
import re
import sqlite3
import subprocess
Expand All @@ -33,6 +34,10 @@
logger = logging.getLogger(__name__)

_DEFAULT_PORT = 9900
# Outcomes that end the request thread's wait but not the task (see A2AAdapter._detach).
_TIMED_OUT = (protocol.STATE_FAILED, "[agent did not reply in time]")
_CLIENT_GONE = (protocol.STATE_FAILED, "[client disconnected]")
_STILL_WORKING = "[still working: poll tasks/get for the reply]"
# 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
Expand Down Expand Up @@ -121,6 +126,19 @@ def _profile_home(profile: str) -> Optional[str]:
return None


class _ExclusiveHTTPServer(ThreadingHTTPServer):
"""A port belongs to one gateway. stdlib sets SO_REUSEADDR, which on Windows lets a second profile bind the
same listening port, and requests then split between them; there the second bind must fail instead
(SO_EXCLUSIVEADDRUSE). Elsewhere SO_REUSEADDR only lets a restart reuse a TIME_WAIT port, so it stays."""

allow_reuse_address = os.name != "nt"

def server_bind(self):
if os.name == "nt" and hasattr(socket, "SO_EXCLUSIVEADDRUSE"):
self.socket.setsockopt(socket.SOL_SOCKET, socket.SO_EXCLUSIVEADDRUSE, 1)
super().server_bind()


def _daemon_thread(target, name: str) -> threading.Thread:
t = threading.Thread(target=target, name=name, daemon=True)
t.start()
Expand Down Expand Up @@ -317,7 +335,7 @@ async def connect(self, **_kwargs) -> bool:
# Capture the gateway loop so the HTTP thread can marshal events via run_coroutine_threadsafe.
self._loop = asyncio.get_running_loop()
try:
self._httpd = ThreadingHTTPServer((self.host, self.port), A2ARequestHandler)
self._httpd = _ExclusiveHTTPServer((self.host, self.port), A2ARequestHandler)
except OSError as e:
logger.error("A2A: could not bind %s:%s β€” %s", self.host, self.port, e)
self._set_fatal_error("bind_failed", f"A2A bind failed: {e}", retryable=True)
Expand Down Expand Up @@ -652,18 +670,33 @@ def _await_future(fut: Future, deadline: float, keepalive, on_timeout: tuple[str
try:
keepalive()
except Exception:
return (protocol.STATE_FAILED, "[client disconnected]")
return _CLIENT_GONE
except Exception:
return on_timeout

def _await_reply(self, pending: dict, keepalive=None) -> tuple[str, str]:
return self._await_future(pending["future"], pending["started"] + _reply_timeout(), keepalive,
(protocol.STATE_FAILED, "[agent did not reply in time]"))
return self._await_future(pending["future"], pending["started"] + _reply_timeout(), keepalive, _TIMED_OUT)

def _detach(self, pending: dict) -> tuple[str, str]:
"""Stop waiting in the request thread but keep the task: it stays WORKING and a daemon thread records the
reply whenever the agent sends it (bounded by the orphan ceiling), so tasks/get, tasks/list and push see
it. A reply past A2A_REPLY_TIMEOUT, or after a stream client left, used to be dropped with the task FAILED."""
deadline = pending["started"] + _MAX_ORPHAN_TIMEOUT
_daemon_thread(lambda: self._finalize_task(pending, *self._await_future(pending["future"], deadline, None, _TIMED_OUT)),
"a2a-task-" + pending["task_id"])
return protocol.STATE_WORKING, _STILL_WORKING

def _finish(self, pending: dict, keepalive=None) -> tuple[str, str]:
"""Wait up to A2A_REPLY_TIMEOUT; past it, or once the stream client is gone, detach instead of failing."""
out = self._await_reply(pending, keepalive)
if out is _TIMED_OUT or out is _CLIENT_GONE:
return self._detach(pending)
return self._finalize_task(pending, *out)

def _rpc_message_send(self, req_id: Any, params: dict, peer: str, agent: Optional[dict] = None, v1_response: bool = False) -> dict:
task, pending = self._prepare_task(params, peer, agent=agent)
if task is None:
state, reply = self._finalize_task(pending, *self._await_reply(pending))
state, reply = self._finish(pending)
task = protocol.build_task(pending["task_id"], pending["context_id"], state, reply, created_at=pending["created_iso"])
return _ok(req_id, protocol.send_message_response(task) if v1_response else task)

Expand Down Expand Up @@ -708,12 +741,12 @@ def _rpc_message_stream(self, handler, req_id: Any, params: dict, peer: str, age
submitted = protocol.build_task(task_id, context_id, protocol.STATE_SUBMITTED, created_at=pending["created_iso"])
self._sse_write(handler, protocol.sse_data(protocol.stream_task(submitted), req_id))
self._sse_write(handler, protocol.sse_data(protocol.status_update(task_id, context_id, protocol.STATE_WORKING), req_id))
state, reply = self._finalize_task(pending, *self._await_reply(pending, keepalive=self._keepalive(handler)))
pending = None
state, reply = self._finish(pending, keepalive=self._keepalive(handler))
pending = None # a detached task ends the stream on WORKING; the client follows it via tasks/get
self._emit_terminal(handler, task_id, context_id, state, reply, req_id=req_id)
except (BrokenPipeError, ConnectionResetError):
if pending is not None:
self._finalize_task(pending, protocol.STATE_FAILED, "[client disconnected]")
self._detach(pending)
logger.debug("A2A: stream client disconnected")

def _rpc_tasks_subscribe(self, handler, req_id: Any, params: dict, agent: Optional[dict] = None) -> None:
Expand Down
16 changes: 12 additions & 4 deletions tests/plugins/test_a2a_phase23.py
Original file line number Diff line number Diff line change
Expand Up @@ -545,7 +545,9 @@ def pause_while_finalizing(reply):
assert result == [(protocol.STATE_COMPLETED, "reply")]
assert adapter.tasks.get("t-live")["state"] == protocol.STATE_COMPLETED

def test_stream_disconnect_releases_active_request(self, monkeypatch):
def test_stream_disconnect_detaches_then_releases_on_the_reply(self, monkeypatch):
"""A stream client that leaves no longer fails the task: the agent's later reply still lands (the
client can come back through tasks/get), and the active request is released then."""
adapter, _base = _make_live_adapter(monkeypatch)
rec = adapter.tasks.create("t-live", "c1", "peer")
adapter.tasks.set_state("t-live", protocol.STATE_WORKING)
Expand Down Expand Up @@ -574,9 +576,15 @@ def end_headers(self):

adapter._rpc_message_stream(Handler(), 1, {}, "peer")

stored = adapter.tasks.get("t-live")
assert stored["state"] == protocol.STATE_FAILED
assert stored["reply"] == "[client disconnected]"
assert adapter.tasks.get("t-live")["state"] == protocol.STATE_WORKING
adapter._resolve_task("t-live", protocol.STATE_COMPLETED, "late reply")
for _ in range(100):
stored = adapter.tasks.get("t-live")
if stored["state"] != protocol.STATE_WORKING:
break
time.sleep(0.02)
assert stored["state"] == protocol.STATE_COMPLETED
assert stored["reply"] == "late reply"
assert "t-live" not in adapter._pending
assert "t-live" not in adapter._active_tasks

Expand Down
66 changes: 57 additions & 9 deletions tests/plugins/test_a2a_plugin.py
Original file line number Diff line number Diff line change
Expand Up @@ -1075,26 +1075,60 @@ async def run():

asyncio.run(run())

def test_timeout_returns_failed_not_completed(self, monkeypatch):
"""When the agent never replies, the task must FAIL (and count as a
failure), not report success."""
def test_timeout_detaches_the_task_and_keeps_the_late_reply(self, monkeypatch):
"""Past A2A_REPLY_TIMEOUT the task is neither failed nor reported as success: it stays WORKING,
and the reply that comes later lands in the task store (it used to be dropped)."""
monkeypatch.delenv("A2A_BEARER_TOKEN", raising=False)
monkeypatch.delenv("A2A_PEER_TOKENS", raising=False)
monkeypatch.setenv("A2A_REPLY_TIMEOUT", "1")
adapter, base = _make_live_adapter(monkeypatch, reply_fn=lambda e: None)

def get_task(task_id):
return _post_json(base + "/", {"jsonrpc": "2.0", "id": "2", "method": "GetTask", "params": {"id": task_id}})["result"]

async def run():
assert await adapter.connect() is True
failed_before = protocol.metrics.tasks_failed
completed_before = protocol.metrics.tasks_completed
resp = await asyncio.to_thread(_post_json, base + "/", _send_body("are you there"))
resp = await asyncio.to_thread(_post_json, base + "/", _send_body("are you there", ctx="ctx-late"))
task = resp["result"]
assert task["status"]["state"] == "TASK_STATE_FAILED"
assert protocol.metrics.tasks_failed == failed_before + 1
assert task["status"]["state"] == protocol.STATE_WORKING
assert "artifacts" not in task
assert protocol.metrics.tasks_failed == failed_before
assert protocol.metrics.tasks_completed == completed_before
# The task store agrees.
rec = adapter.tasks.get(task["id"])
assert rec["state"] == "TASK_STATE_FAILED"
await adapter.send("ctx-late", "late answer", metadata={"notify": True})
for _ in range(50):
done = await asyncio.to_thread(get_task, task["id"])
if done["status"]["state"] != protocol.STATE_WORKING:
break
await asyncio.sleep(0.1)
assert done["status"]["state"] == protocol.STATE_COMPLETED
assert protocol.extract_text(done["artifacts"][0]) == "late answer"
await adapter.disconnect()

asyncio.run(run())

def test_a_detached_task_that_never_gets_a_reply_fails_at_the_ceiling(self, monkeypatch):
"""The detached wait is bounded: no reply by the orphan ceiling => FAILED, counted as a failure."""
import plugins.platforms.a2a.adapter as adapter_mod
monkeypatch.delenv("A2A_BEARER_TOKEN", raising=False)
monkeypatch.delenv("A2A_PEER_TOKENS", raising=False)
monkeypatch.setenv("A2A_REPLY_TIMEOUT", "1")
monkeypatch.setattr(adapter_mod, "_MAX_ORPHAN_TIMEOUT", 2)
adapter, base = _make_live_adapter(monkeypatch, reply_fn=lambda e: None)

async def run():
assert await adapter.connect() is True
failed_before = protocol.metrics.tasks_failed
task = (await asyncio.to_thread(_post_json, base + "/", _send_body("silence")))["result"]
assert task["status"]["state"] == protocol.STATE_WORKING
for _ in range(60):
rec = adapter.tasks.get(task["id"])
if rec["state"] != protocol.STATE_WORKING:
break
await asyncio.sleep(0.1)
assert rec["state"] == protocol.STATE_FAILED
assert protocol.metrics.tasks_failed == failed_before + 1
await adapter.disconnect()

asyncio.run(run())
Expand Down Expand Up @@ -1613,6 +1647,20 @@ def test_forward_to_profile_first_contact_creates_then_resumes_fake_hermes(self,
assert title == "a2a-dev-ctx-unsafe-value"


class TestExclusivePort:
"""Two gateways (profiles) must not share one A2A port: on Windows stdlib SO_REUSEADDR let a second
server bind the same listening port and requests split between them."""

def test_a_second_server_on_the_same_port_fails_to_bind(self):
from plugins.platforms.a2a.adapter import A2ARequestHandler, _ExclusiveHTTPServer
first = _ExclusiveHTTPServer(("127.0.0.1", 0), A2ARequestHandler)
try:
with pytest.raises(OSError):
_ExclusiveHTTPServer(("127.0.0.1", first.server_address[1]), A2ARequestHandler).server_close()
finally:
first.server_close()


# --------------------------------------------------------------------------
# Multiplex secondary-profile scope (construction-time config leak)
# --------------------------------------------------------------------------
Expand Down
4 changes: 2 additions & 2 deletions website/docs/user-guide/messaging/a2a.md
Original file line number Diff line number Diff line change
Expand Up @@ -102,7 +102,7 @@ Secure by default; every widening step is explicit:
| `A2A_ALLOW_ALL_USERS` | `false` | Allow any authenticated peer (dev only) |
| `A2A_RATE_LIMIT` | `60` | Requests/minute per identity |
| `A2A_MAX_PINGPONG_TURNS` | `5` | Anti-loop turn cap per context (max 20) |
| `A2A_REPLY_TIMEOUT` | `300` | Seconds to wait for the agent's reply. The orphan-task sweep never fails a task before this window elapses (floor 300s), and never while a request is still waiting on it |
| `A2A_REPLY_TIMEOUT` | `300` | Seconds a request waits for the agent's reply. Past it the task is not failed: the caller gets it in the working state, and the reply is stored when it arrives (a task that never gets one fails at the orphan ceiling, 24 h) |
| `A2A_PUSH_SECRET` | bearer token | HMAC secret for push-notification signing |
| `A2A_ADVERTISED_TOOLSETS` | all registered | Restrict which skills appear on the Agent Card |

Expand All @@ -127,4 +127,4 @@ curl -X POST http://your-host:9900/ \
- **Peers can't reach the card URL** β€” the card was advertising your bind address; set `A2A_PUBLIC_URL` to the externally routable URL.
- **`401 Unauthorized`** β€” token mismatch; check `A2A_PEER_TOKENS`/`A2A_BEARER_TOKEN` on the server and the peer's `auth:` block.
- **Server won't bind non-localhost** β€” by design: set a bearer token first, then `A2A_HOST=0.0.0.0`.
- **Replies time out on long tasks** β€” raise `A2A_REPLY_TIMEOUT` (the orphan sweep follows it, so a late reply is stored, not discarded), or have the caller register a push-notification config and poll `GetTask`.
- **Replies take longer than `A2A_REPLY_TIMEOUT`** β€” the caller gets the task in the working state and fetches the reply with `GetTask` (or registers a push-notification config); the same holds when a streaming client disconnects. Raise `A2A_REPLY_TIMEOUT` to keep blocking callers waiting longer.