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
36 changes: 30 additions & 6 deletions plugins/platforms/a2a/DESIGN.md
Original file line number Diff line number Diff line change
Expand Up @@ -19,10 +19,13 @@ must not touch core files.** A2A now lives entirely under
### Outbound — client tools (`a2a` toolset)
- `a2a_discover(url)` — fetch + summarize a peer's Agent Card (v1.0
`supportedInterfaces` aware, tolerates 0.3 cards).
- `a2a_call(agent, message, context_id?)` — send a JSON-RPC `message/send`
- `a2a_call(agent, message, context_id?, return_immediately?)` — send a JSON-RPC `message/send`
task to a peer, return the reply. Multi-turn via `context_id` (carried
inside the Message per v1.0). Surfaces `TASK_STATE_INPUT_REQUIRED` so the
model knows to answer and continue the context.
model knows to answer and continue the context. `return_immediately` sets
`configuration.returnImmediately` so the peer returns a working task id
instead of holding the HTTP call until the job finishes.
- `a2a_get_task(agent, task_id)` — JSON-RPC `GetTask` poll for that id.
- `a2a_list()` — configured peers + persisted conversations + metrics.
- `a2a_history(context_id, limit?)` — recall a persisted conversation
(this is the production consumer of the persistence layer).
Expand All @@ -47,18 +50,29 @@ Peers resolved from `config.yaml` → `a2a_agents`, or a direct URL.
- JSON-RPC methods: `message/send`, `message/stream` (SSE), `tasks/get`,
`tasks/list`, `tasks/cancel`, `tasks/subscribe`,
`tasks/pushNotificationConfig/create` (legacy `set` names accepted).
`message/send` honors `configuration.returnImmediately` (and the older
`configuration.blocking: false`): the HTTP call returns a working Task at
once, a background waiter records the real result, and `tasks/get` can
poll it. Blocking send is unchanged (wait until done or `A2A_REPLY_TIMEOUT`).
The orphan watchdog skips task ids that still have a live waiter, so a
long job is not marked failed just because it has been working for more
than five minutes. The waiter itself stops after 24 hours
(`_BACKGROUND_WAIT_SECONDS`) and records the task as failed, so a hung
gateway turn cannot pin a daemon thread forever.
- **Live-session injection (the #11025 insight):** inbound tasks route through
the normal `MessageEvent` → `handle_message` path keyed by the A2A
`contextId`, so the agent that answers is the same one serving the user —
full memory/context, not a clone. The reply returns through `adapter.send()`,
which fulfils the pending per-**task** `Future` the HTTP request is blocked
on (per-context FIFO, so concurrent same-context requests can't cross-talk);
on. The final's reply anchor (the task id) picks the task, and the adapter
dispatches one turn per context at a time, so the gateway busy queue never
folds two same-context tasks into one turn;
`on_processing_complete` resolves failures/cancellations promptly.
- **Task store:** every task (including terminal ones, bounded to the last
500) stays queryable via `tasks/get` / `tasks/list`, and `tasks/subscribe`
reattaches to a running task's stream via store watchers. A watchdog fails
orphaned tasks after 5 minutes (idempotent transitions — no double
counting in metrics).
orphaned tasks after 5 minutes if they have no live waiter (idempotent
transitions, no double counting in metrics).
- **input-required:** the platform hint tells the agent to start a reply with
`[INPUT_REQUIRED]` when it needs clarification; the adapter maps that to
`TASK_STATE_INPUT_REQUIRED` with the question in `status.message`.
Expand Down Expand Up @@ -124,9 +138,19 @@ outbound client tools (`/metrics` and `a2a_list` report both directions).

## Persistence (survives compaction)
A2A conversations are written to `~/.hermes/a2a_conversations/<context>.jsonl`,
outside the context-compaction pipeline — compaction and restarts can't lose
outside the context-compaction pipeline. Compaction and restarts can't lose
them (#11025 requirement). The `a2a_history` tool recalls them by context id.

Outbound `a2a_call` with `return_immediately` writes the prompt at send time.
The peer reply is written only when `a2a_get_task` sees a terminal state. An
unpolled job therefore has a user line and no agent line in `a2a_history`.
That is a known property, not a missing persist. Inbound tasks still write
both sides, because the background waiter records the reply itself.

`a2a_call` and `a2a_get_task` cache the peer Agent Card for 60 seconds so a
polling loop does not refetch it on every GetTask. `a2a_discover` always
fetches a fresh card.

## Requirements traced to the cluster

| Source | Requirement | Where |
Expand Down
5 changes: 3 additions & 2 deletions plugins/platforms/a2a/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -31,10 +31,11 @@ a2a_agents:

## Outbound — call other agents

The agent gets five tools:
The agent gets these tools:

- `a2a_discover(url)` — what can this agent do?
- `a2a_call(agent, message, context_id?)` — send it a task, get the reply.
- `a2a_call(agent, message, context_id?, return_immediately?)` — send it a task, get the reply. Set `return_immediately` for long jobs to get a task id back at once.
- `a2a_get_task(agent, task_id)` — check a running task (state and, when done, the reply).
- `a2a_list()` — configured peers, saved conversations, metrics.
- `a2a_history(context_id)` — recall a saved A2A conversation.
- `a2a_orchestrate(capability, message, mode?)` — fan-out a task to every
Expand Down
151 changes: 133 additions & 18 deletions plugins/platforms/a2a/adapter.py
Original file line number Diff line number Diff line change
Expand Up @@ -18,6 +18,7 @@
import urllib.parse
import urllib.request
from collections import deque
from contextvars import copy_context
from concurrent.futures import Future
from concurrent.futures import TimeoutError as FuturesTimeout
from http.server import BaseHTTPRequestHandler, ThreadingHTTPServer
Expand All @@ -32,6 +33,7 @@

logger = logging.getLogger(__name__)

_BACKGROUND_WAIT_SECONDS = 24 * 60 * 60
_DEFAULT_PORT = 9900
# 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).
Expand Down Expand Up @@ -291,10 +293,14 @@ def __init__(self, config, **kwargs):
self._profile_sessions: Dict[tuple[str, str, str], str] = {}
self._profile_session_locks: Dict[tuple[str, str, str], threading.Lock] = {}
self._profile_session_locks_guard = threading.Lock()
# Pending reply futures: task_id -> (context_id, Future). _pending_order keeps per-context
# FIFO so adapter.send() — which only knows the context — resolves the oldest task.
# Pending reply futures: task_id -> (context_id, Future). send() resolves the task named by
# the final's reply anchor (the gateway anchors on the inbound message id == task id).
self._pending: Dict[str, tuple[str, Future]] = {}
self._pending_order: Dict[str, deque[str]] = {}
# One dispatched turn per context. The gateway's busy queue merges or replaces queued text
# for a session, which would fold two tasks into one turn: one gets the other's answer and
# the other never settles. Later tasks wait here until the in-flight one is popped.
self._inflight: Dict[str, str] = {}
self._queued: Dict[str, deque[tuple[str, MessageEvent]]] = {}
# Request ownership outlives reply Futures and also covers synchronous profile forwards.
self._active_tasks: set[str] = set()
self._pending_lock = threading.Lock()
Expand Down Expand Up @@ -346,7 +352,8 @@ async def disconnect(self) -> None:
for tid in list(self._pending):
self._resolve_locked(tid, protocol.STATE_FAILED, "[agent shutting down]")
self._pending.clear()
self._pending_order.clear()
self._inflight.clear()
self._queued.clear()
self._active_tasks.clear()

def _watchdog_loop(self) -> None:
Expand Down Expand Up @@ -480,7 +487,6 @@ def _add_pending(self, task_id: str, context_id: str) -> Future:
with self._pending_lock:
self._active_tasks.add(task_id)
self._pending[task_id] = (context_id, fut)
self._pending_order.setdefault(context_id, deque()).append(task_id)
return fut

def _activate_task(self, task_id: str) -> None:
Expand All @@ -491,11 +497,38 @@ def _pop_pending(self, task_id: str) -> None:
with self._pending_lock:
self._active_tasks.discard(task_id)
entry = self._pending.pop(task_id, None)
order = self._pending_order.get(entry[0]) if entry else None
if order and task_id in order:
order.remove(task_id)
if order is not None and not order:
self._pending_order.pop(entry[0], None)
nxt = self._advance_context_locked(entry[0]) if entry and self._inflight.get(entry[0]) == task_id else None
if nxt is not None:
self._dispatch_queued(*nxt)

def _claim_context(self, context_id: str, task_id: str, event: MessageEvent) -> bool:
"""True if ``task_id`` may dispatch now; otherwise it waits behind the in-flight task."""
with self._pending_lock:
if context_id in self._inflight:
self._queued.setdefault(context_id, deque()).append((task_id, event))
return False
self._inflight[context_id] = task_id
return True

def _advance_context_locked(self, context_id: str) -> Optional[tuple[str, MessageEvent]]:
"""Hand the context to the next queued task still waiting (cancelled ones were popped)."""
queue = self._queued.get(context_id)
while queue:
task_id, event = queue.popleft()
entry = self._pending.get(task_id)
if entry and not entry[1].done():
self._inflight[context_id] = task_id
return task_id, event
self._queued.pop(context_id, None)
self._inflight.pop(context_id, None)
return None

def _dispatch_queued(self, task_id: str, event: MessageEvent) -> None:
try:
asyncio.run_coroutine_threadsafe(self.handle_message(event), self._loop)
except Exception as e:
# Its waiter finalizes the failure and pops it, which hands the context on.
self._resolve_task(task_id, protocol.STATE_FAILED, security.redact_outbound(f"Dispatch failed: {e}"))

def _resolve_locked(self, task_id: str, state: str, text: str) -> bool:
entry = self._pending.get(task_id)
Expand All @@ -508,9 +541,14 @@ def _resolve_task(self, task_id: str, state: str, text: str) -> bool:
with self._pending_lock:
return self._resolve_locked(task_id, state, text)

def _resolve_oldest_for_context(self, context_id: str, state: str, text: str) -> bool:
def _resolve_final(self, context_id: str, anchor: Optional[str], text: str) -> bool:
"""Resolve the task that owns a final: its anchor, else the context's in-flight task. A
late final anchored on a finished task settles nothing, so it can't answer a sibling."""
with self._pending_lock:
return any(self._resolve_locked(tid, state, text) for tid in self._pending_order.get(context_id, ()))
task_id = anchor or self._inflight.get(context_id)
entry = self._pending.get(task_id or "")
return bool(entry and entry[0] == context_id
and self._resolve_locked(task_id, protocol.STATE_COMPLETED, text))

def _scope_for_agent(self, agent: Optional[dict]) -> tuple[str, str]:
return tuple(str((agent or self._agents[""]).get(k) or "") for k in ("slug", "tenant"))
Expand All @@ -525,7 +563,7 @@ def _end_task(self, rec: dict, state: str, text: str, stored_reply: str = "") ->
protocol.metrics.tasks_failed += state == protocol.STATE_FAILED
return protocol.build_task(rec["task_id"], rec["context_id"], state, text, created_at=rec["created_iso"]), None

def _prepare_task(self, params: dict, peer: str, agent: Optional[dict] = None) -> tuple[Optional[dict], Optional[dict]]:
def _prepare_task(self, params: dict, peer: str, agent: Optional[dict] = None, *, defer_forward: bool = False) -> tuple[Optional[dict], Optional[dict]]:
"""Validate, register, and dispatch an inbound message (HTTP worker thread). Returns
(terminal_task, None) when it ends immediately, else (None, pending) with the future to wait on."""
agent = agent or self._agents[""]
Expand All @@ -549,6 +587,11 @@ def _prepare_task(self, params: dict, peer: str, agent: Optional[dict] = None) -
self._register_inline_push(task_id, params, agent=agent)
if not agent.get("local", True):
self._activate_task(task_id)
if defer_forward:
self.tasks.set_state(task_id, protocol.STATE_WORKING)
return None, {"task_id": task_id, "context_id": context_id, "peer": peer,
"created_iso": rec["created_iso"], "started": time.time(),
"forward": ({**agent, "timeout": _BACKGROUND_WAIT_SECONDS}, peer, context_id, framed)}
try:
reply, state = self._forward_to_profile(agent, peer, context_id, framed)
self._record_outcome(task_id, context_id, peer, state, reply)
Expand All @@ -561,7 +604,8 @@ def _prepare_task(self, params: dict, peer: str, agent: Optional[dict] = None) -
event = MessageEvent(text=framed, message_type=MessageType.TEXT, message_id=task_id,
source=self.build_source(chat_id=context_id, chat_name=f"a2a:{peer}", chat_type="dm", user_id=peer, user_name=peer))
try:
asyncio.run_coroutine_threadsafe(self.handle_message(event), self._loop)
if self._claim_context(context_id, task_id, event):
asyncio.run_coroutine_threadsafe(self.handle_message(event), self._loop)
except Exception as e:
msg = security.redact_outbound(f"Dispatch failed: {e}")
try:
Expand Down Expand Up @@ -628,6 +672,10 @@ def _finalize_task(self, pending: dict, state: str, reply: str) -> tuple[str, st
"""Record a dispatched task's outcome; returns (state, reply) after redaction and
input-required detection (a leading marker flags a clarification request)."""
task_id, context_id, peer = pending["task_id"], pending["context_id"], pending["peer"]
rec = self.tasks.get(task_id)
if rec and rec["state"] in protocol.TERMINAL_STATES:
self._pop_pending(task_id)
return rec["state"], rec.get("reply") or ""
try:
reply = security.redact_outbound(reply or "")
stripped = reply.lstrip()
Expand Down Expand Up @@ -660,8 +708,73 @@ 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]"))

@staticmethod
def _config_return_immediately(params: dict) -> bool:
"""True when the caller asked for a task id now, not a blocking wait.

A2A v1.0 uses configuration.returnImmediately. Older peers used
configuration.blocking = false. Either one is enough.
"""
cfg = params.get("configuration") if isinstance(params, dict) else None
if not isinstance(cfg, dict):
return False
if cfg.get("returnImmediately") is True or cfg.get("return_immediately") is True:
return True
if cfg.get("blocking") is False:
return True
return False

def _wait_in_background(self, pending: dict) -> None:
"""Wait for the agent, then record the result.

Used when returnImmediately is set. The HTTP caller already has the
working task. A2A_REPLY_TIMEOUT applies only to blocking callers.
The wait still stops after _BACKGROUND_WAIT_SECONDS so a hung
gateway turn cannot pin a daemon thread and a forever-WORKING task.
"""

def _run() -> None:
try:
try:
if "forward" in pending:
reply, state = self._forward_to_profile(*pending["forward"])
else:
state, reply = pending["future"].result(timeout=_BACKGROUND_WAIT_SECONDS)
except FuturesTimeout:
logger.warning(
"A2A: background waiter for task %s hit the %ss ceiling",
pending.get("task_id"),
_BACKGROUND_WAIT_SECONDS,
)
state, reply = (
protocol.STATE_FAILED,
"[agent did not reply in time]",
)
except Exception:
state, reply = protocol.STATE_FAILED, "[agent did not reply]"
self._finalize_task(pending, state, reply)
except Exception:
logger.debug("A2A: background waiter failed", exc_info=True)
try:
self._finalize_task(
pending, protocol.STATE_FAILED, "[agent did not reply]")
except Exception:
pass

threading.Thread(
target=copy_context().run,
args=(_run,),
name=f"a2a-wait-{pending['task_id'][:12]}",
daemon=True,
).start()

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)
options = {"defer_forward": True} if agent and not agent.get("local", True) and self._config_return_immediately(params) else {}
task, pending = self._prepare_task(params, peer, agent=agent, **options)
if task is None and self._config_return_immediately(params):
self._wait_in_background(pending)
task = protocol.build_task(pending["task_id"], pending["context_id"], protocol.STATE_WORKING,
created_at=pending["created_iso"])
if task is None:
state, reply = self._finalize_task(pending, *self._await_reply(pending))
task = protocol.build_task(pending["task_id"], pending["context_id"], state, reply, created_at=pending["created_iso"])
Expand Down Expand Up @@ -826,13 +939,15 @@ def fail(msg: str, *args) -> None:
logger.debug("A2A: push notification sent for task %s", task_id)

async def send(self, chat_id: str, content: str, reply_to: Optional[str] = None, metadata: Optional[Dict[str, Any]] = None):
"""Fulfil the oldest pending reply Future for this context (``chat_id`` = A2A context id).
"""Fulfil the reply Future of the task this final belongs to (``chat_id`` = A2A context id).
Only sends carrying ``metadata['notify']`` (the base adapter's final-reply marker) satisfy
the caller; progress/status/preview sends must not."""
# Stream-consumer finals carry their anchor in metadata instead of ``reply_to``.
anchor = str(reply_to or (metadata or {}).get("reply_to_message_id") or "") or 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
elif not self._resolve_final(chat_id, anchor, content or ""):
logger.info("A2A: final for context %s (anchor %s) matched no pending task", chat_id, anchor)
return SendResult(success=True, message_id=str(int(time.time() * 1000)))

async def send_typing(self, chat_id: str, metadata=None) -> None:
Expand Down
Loading