Skip to content
Closed
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
23 changes: 20 additions & 3 deletions plugins/platforms/a2a/adapter.py
Original file line number Diff line number Diff line change
Expand Up @@ -657,8 +657,24 @@ def _await_future(fut: Future, deadline: float, keepalive, on_timeout: tuple[str
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]"))
"""The observation deadline is not a processing failure while the task is live."""
deadline = pending["started"] + _reply_timeout()
logged_timeout = False
while True:
try:
return pending["future"].result(timeout=_SSE_KEEPALIVE)
except FuturesTimeout:
if time.time() >= deadline and not logged_timeout:
logger.info("A2A: observation window elapsed; preserving reply waiter task=%s context=%s",
pending["task_id"], pending["context_id"])
logged_timeout = True
if keepalive:
try:
keepalive()
except Exception:
return protocol.STATE_FAILED, "[client disconnected]"
except Exception:
return protocol.STATE_FAILED, "[agent processing failed]"

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)
Expand Down Expand Up @@ -832,7 +848,8 @@ async def send(self, chat_id: str, content: str, reply_to: Optional[str] = None,
if not (metadata or {}).get("notify"):
logger.debug("A2A: ignoring non-final send for context %s", chat_id)
elif not self._resolve_oldest_for_context(chat_id, protocol.STATE_COMPLETED, content or ""):
logger.debug("A2A: send() for context %s had no pending waiter", chat_id) # late chunk / out-of-band
logger.error("A2A: final reply undelivered: context=%s task=%s: no pending waiter", chat_id, reply_to or "unknown")
return SendResult(success=False, error="no pending waiter; final reply undelivered")
return SendResult(success=True, message_id=str(int(time.time() * 1000)))

async def send_typing(self, chat_id: str, metadata=None) -> None:
Expand Down
36 changes: 36 additions & 0 deletions tests/plugins/test_a2a_late_reply.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,36 @@
"""Final reply delivery must never claim success without a recipient."""
import asyncio
import logging

from gateway.config import PlatformConfig
from plugins.platforms.a2a.adapter import A2AAdapter


def test_final_without_waiter_fails_loudly(caplog):
adapter = A2AAdapter(PlatformConfig(enabled=True))
with caplog.at_level(logging.WARNING):
result = asyncio.run(adapter.send("ctx-late", "late answer", reply_to="task-expired", metadata={"notify": True}))
assert result.success is False
assert "no pending waiter" in result.error
assert "ctx-late" in caplog.text
assert "task-expired" in caplog.text


def test_reply_after_observation_deadline_reaches_original_task(monkeypatch):
import time
from plugins.platforms.a2a import adapter as module, protocol
monkeypatch.setattr(module, "_SSE_KEEPALIVE", 0.001)
adapter = A2AAdapter(PlatformConfig(enabled=True))
rec = adapter.tasks.create("task-late", "ctx-late", "caller")
future = adapter._add_pending("task-late", "ctx-late")
pending = {"task_id": "task-late", "context_id": "ctx-late", "peer": "caller",
"future": future, "started": time.time() - 1000, "created_iso": rec["created_iso"]}

def late_reply():
result = asyncio.run(adapter.send("ctx-late", "late answer", reply_to="task-late", metadata={"notify": True}))
assert result.success

state, text = adapter._finalize_task(pending, *adapter._await_reply(pending, keepalive=late_reply))
assert state == protocol.STATE_COMPLETED
assert text == "late answer"
assert adapter.tasks.get("task-late")["reply"] == "late answer"