feat(a2a): implement A2A stdio transport for local mode
CI / push-validation (pull_request) Successful in 32s
CI / helm (pull_request) Successful in 38s
CI / build (pull_request) Successful in 4m5s
CI / lint (pull_request) Successful in 4m20s
CI / quality (pull_request) Successful in 4m39s
CI / typecheck (pull_request) Successful in 4m52s
CI / security (pull_request) Successful in 5m3s
CI / e2e_tests (pull_request) Successful in 7m3s
CI / integration_tests (pull_request) Successful in 7m38s
CI / unit_tests (pull_request) Successful in 9m30s
CI / docker (pull_request) Successful in 1m38s
CI / coverage (pull_request) Successful in 15m31s
CI / benchmark-regression (push) Waiting to run
CI / benchmark-publish (push) Waiting to run
CI / status-check (pull_request) Successful in 3s
CI / push-validation (push) Successful in 22s
CI / helm (push) Successful in 30s
CI / build (push) Successful in 3m49s
CI / lint (push) Successful in 3m57s
CI / quality (push) Successful in 4m23s
CI / typecheck (push) Successful in 4m41s
CI / security (push) Successful in 4m48s
CI / unit_tests (push) Successful in 7m21s
CI / integration_tests (push) Successful in 7m28s
CI / e2e_tests (push) Successful in 7m34s
CI / docker (push) Successful in 1m38s
CI / coverage (push) Successful in 15m20s
CI / status-check (push) Successful in 3s
CI / benchmark-publish (pull_request) Has been skipped
CI / benchmark-regression (pull_request) Failing after 1h1m27s

This commit was merged in pull request #10791.
This commit is contained in:
2026-04-20 06:57:37 +00:00
committed by Forgejo
parent 58fa90b1b5
commit b3bfbc1d55
3 changed files with 308 additions and 1 deletions
+9 -1
View File
@@ -5,7 +5,8 @@ models, event streaming stubs, and version negotiation for the A2A boundary.
In **local mode** the :class:`A2aLocalFacade` maps A2A operation names to In **local mode** the :class:`A2aLocalFacade` maps A2A operation names to
direct Python method calls on existing application services. No serialization, direct Python method calls on existing application services. No serialization,
no network, no authentication. no network, no authentication. The :class:`A2aStdioTransport` handles
subprocess communication via JSON-RPC 2.0 over stdin/stdout.
In **server mode** the :class:`A2aHttpTransport` is a stub that raises In **server mode** the :class:`A2aHttpTransport` is a stub that raises
:class:`A2aNotAvailableError` for every operation. When server mode is :class:`A2aNotAvailableError` for every operation. When server mode is
@@ -14,6 +15,9 @@ implemented the concrete transport will replace these stubs.
Server client protocols (:class:`ServerClient`, :class:`RemoteExecutionClient`, Server client protocols (:class:`ServerClient`, :class:`RemoteExecutionClient`,
:class:`AuthClient`) and their stub implementations are provided for forward :class:`AuthClient`) and their stub implementations are provided for forward
compatibility. :class:`ServerConnectionConfig` validates connection parameters. compatibility. :class:`ServerConnectionConfig` validates connection parameters.
The :class:`TransportSelector` chooses the appropriate transport based on
configuration: stdio for local mode, HTTP for server mode.
""" """
from __future__ import annotations from __future__ import annotations
@@ -42,7 +46,9 @@ from cleveragents.a2a.models import (
A2aVersion, A2aVersion,
) )
from cleveragents.a2a.server_config import ServerConnectionConfig from cleveragents.a2a.server_config import ServerConnectionConfig
from cleveragents.a2a.stdio_transport import A2aStdioTransport
from cleveragents.a2a.transport import A2aHttpTransport from cleveragents.a2a.transport import A2aHttpTransport
from cleveragents.a2a.transport_selector import TransportSelector
from cleveragents.a2a.versioning import A2aVersionNegotiator from cleveragents.a2a.versioning import A2aVersionNegotiator
__all__ = [ __all__ = [
@@ -56,6 +62,7 @@ __all__ = [
"A2aOperationNotFoundError", "A2aOperationNotFoundError",
"A2aRequest", "A2aRequest",
"A2aResponse", "A2aResponse",
"A2aStdioTransport",
"A2aVersion", "A2aVersion",
"A2aVersionMismatchError", "A2aVersionMismatchError",
"A2aVersionNegotiator", "A2aVersionNegotiator",
@@ -66,4 +73,5 @@ __all__ = [
"StubAuthClient", "StubAuthClient",
"StubRemoteExecutionClient", "StubRemoteExecutionClient",
"StubServerClient", "StubServerClient",
"TransportSelector",
] ]
+241
View File
@@ -0,0 +1,241 @@
"""A2A local-mode stdio transport for subprocess communication.
Implements JSON-RPC 2.0 message framing over stdin/stdout for communicating
with an agent subprocess in local mode. The CLI spawns the agent as a
subprocess and sends JSON-RPC requests over stdin, receiving responses
over stdout.
"""
from __future__ import annotations
import json
import subprocess
import sys
import structlog
from cleveragents.a2a.models import A2aRequest, A2aResponse
logger: structlog.stdlib.BoundLogger = structlog.get_logger(__name__)
class A2aStdioTransport:
"""Stdio transport for local-mode subprocess communication.
Manages a subprocess and communicates with it via JSON-RPC 2.0 messages
over stdin/stdout. Each message is a single JSON object followed by
a newline.
"""
def __init__(self) -> None:
"""Initialize the stdio transport."""
self._process: subprocess.Popen[str] | None = None
self._is_connected: bool = False
def send(self, request: A2aRequest) -> A2aResponse:
"""Send an A2A request over stdio and receive the response.
Args:
request: The A2aRequest to send.
Returns:
The A2aResponse received from the subprocess.
Raises:
RuntimeError: If not connected to a subprocess.
ValueError: If request is not an A2aRequest instance.
"""
if not isinstance(request, A2aRequest):
raise TypeError("request must be an A2aRequest instance")
if not self._is_connected or self._process is None:
raise RuntimeError("Not connected to subprocess")
# Serialize request to JSON-RPC 2.0 format
request_dict = request.model_dump(exclude_none=True)
request_json = json.dumps(request_dict)
try:
# Send request over stdin
if self._process.stdin is None:
raise RuntimeError("Subprocess stdin is not available")
self._process.stdin.write(request_json + "\n")
self._process.stdin.flush()
logger.debug(
"a2a.stdio.send",
method=request.method,
request_id=request.id,
)
# Read response from stdout
if self._process.stdout is None:
raise RuntimeError("Subprocess stdout is not available")
response_line = self._process.stdout.readline()
if not response_line:
raise RuntimeError("Subprocess closed unexpectedly")
response_dict = json.loads(response_line.strip())
response = A2aResponse(**response_dict)
logger.debug(
"a2a.stdio.receive",
method=request.method,
request_id=request.id,
has_error=response.error is not None,
)
return response
except json.JSONDecodeError as exc:
logger.error(
"a2a.stdio.json_decode_error",
method=request.method,
request_id=request.id,
error=str(exc),
)
raise RuntimeError(f"Invalid JSON response from subprocess: {exc}") from exc
except Exception as exc:
logger.error(
"a2a.stdio.send_error",
method=request.method,
request_id=request.id,
error=str(exc),
)
raise
def connect(self, agent_path: str, *args: str) -> None:
"""Launch the agent subprocess.
Args:
agent_path: Path to the agent executable or Python module.
*args: Additional arguments to pass to the agent.
Raises:
ValueError: If agent_path is empty or not a string.
RuntimeError: If subprocess launch fails.
"""
if not agent_path or not isinstance(agent_path, str):
raise ValueError("agent_path must be a non-empty string")
if self._is_connected:
raise RuntimeError("Already connected to a subprocess")
try:
# Construct command: python -m cleveragents.a2a.cli_bootstrap [args]
# or direct path to agent executable
if agent_path.endswith(".py") or agent_path.startswith("cleveragents."):
# Python module path
cmd = [sys.executable, "-m", agent_path, *list(args)]
else:
# Direct executable path
cmd = [agent_path, *list(args)]
logger.info(
"a2a.stdio.connect",
agent_path=agent_path,
cmd=" ".join(cmd),
)
self._process = subprocess.Popen(
cmd,
stdin=subprocess.PIPE,
stdout=subprocess.PIPE,
stderr=subprocess.PIPE,
text=True,
bufsize=1, # Line buffering
)
self._is_connected = True
logger.info(
"a2a.stdio.connected",
pid=self._process.pid,
)
except FileNotFoundError as exc:
logger.error(
"a2a.stdio.agent_not_found",
agent_path=agent_path,
error=str(exc),
)
raise RuntimeError(f"Agent not found: {agent_path}") from exc
except Exception as exc:
logger.error(
"a2a.stdio.connect_error",
agent_path=agent_path,
error=str(exc),
)
raise RuntimeError(f"Failed to launch agent: {exc}") from exc
def disconnect(self) -> None:
"""Close the connection to the subprocess.
Terminates the subprocess gracefully, waiting for it to exit.
"""
if not self._is_connected or self._process is None:
return
try:
logger.info(
"a2a.stdio.disconnect",
pid=self._process.pid,
)
# Close stdin to signal EOF to subprocess
if self._process.stdin is not None:
self._process.stdin.close()
# Wait for subprocess to exit gracefully
try:
self._process.wait(timeout=5.0)
except subprocess.TimeoutExpired:
logger.warning(
"a2a.stdio.terminate",
pid=self._process.pid,
)
self._process.terminate()
try:
self._process.wait(timeout=2.0)
except subprocess.TimeoutExpired:
logger.error(
"a2a.stdio.kill",
pid=self._process.pid,
)
self._process.kill()
self._process.wait()
self._is_connected = False
logger.info(
"a2a.stdio.disconnected",
pid=self._process.pid,
)
except Exception as exc:
logger.error(
"a2a.stdio.disconnect_error",
error=str(exc),
)
self._is_connected = False
def is_connected(self) -> bool:
"""Return connection status.
Returns:
``True`` if connected to a subprocess, ``False`` otherwise.
"""
return self._is_connected
def get_process(self) -> subprocess.Popen[str] | None:
"""Return the subprocess handle.
Returns:
The subprocess Popen object, or None if not connected.
"""
return self._process
__all__ = [
"A2aStdioTransport",
]
@@ -0,0 +1,58 @@
"""Transport selector for choosing between stdio and HTTP transports.
Selects the appropriate A2A transport based on configuration:
- Stdio transport for local mode (no server URL configured)
- HTTP transport for server mode (server URL configured)
"""
from __future__ import annotations
from typing import TYPE_CHECKING, Union
import structlog
if TYPE_CHECKING:
from cleveragents.a2a.stdio_transport import A2aStdioTransport
from cleveragents.a2a.transport import A2aHttpTransport
logger: structlog.stdlib.BoundLogger = structlog.get_logger(__name__)
# Type alias for transport union
A2aTransport = Union["A2aStdioTransport", "A2aHttpTransport"]
class TransportSelector:
"""Selects the appropriate A2A transport based on configuration.
In local mode (no server URL), selects the stdio transport.
In server mode (server URL configured), selects the HTTP transport.
"""
@staticmethod
def select(server_url: str | None = None) -> A2aTransport:
"""Select the appropriate transport.
Args:
server_url: The server URL for server mode, or None for local mode.
Returns:
An A2aStdioTransport for local mode, or A2aHttpTransport for server mode.
"""
if not server_url:
# Local mode: use stdio transport
from cleveragents.a2a.stdio_transport import A2aStdioTransport
logger.debug("a2a.transport_selector.selected_stdio")
return A2aStdioTransport()
else:
# Server mode: use HTTP transport
from cleveragents.a2a.transport import A2aHttpTransport
logger.debug("a2a.transport_selector.selected_http", server_url=server_url)
return A2aHttpTransport()
__all__ = [
"A2aTransport",
"TransportSelector",
]