[MCP在LangChain中的应用-05]如何实现基于反向通信的进度报告、日志回传和信息征询
A2A SDK服务端框架旨在针对一种或者多种协议绑定构建一个Agent Server。除了默认支持的三种协议绑定(JSON-RPC、gRPC和HTTP+JSON/REST)之外,A2A SDK服务端框架还具有高度的可扩展性,允许开发者根据需要实现自定义的协议绑定。我们只关注JSON-RPC相关的部分,这部分的涉及基本上体现在如下所示的UML中。
1. 利用JSONRPCAppication构建ASGIApplication
通过上面的实例演示我们知道,由A2A SDK服务端框架构建的A2A Server本质上是一个由Uvicorn构建的Web Server,所以我们需要为它构建一个ASGIApplication来处理由Uvicorn接收的请求并作用相应的响应。这个ASGIApplication就是通过JSONRPCApplication来构建的。如下面的代码片段所示,JSONRPCApplication的抽象方法build会返回一个FastAPI或者Starlette对象,作为ASGIApplication被Uvicorn使用。
class JSONRPCApplication(ABC):
def __init__(
self,
agent_card: AgentCard,
http_handler: RequestHandler,
extended_agent_card: AgentCard | None = None,
context_builder: CallContextBuilder | None = None,
card_modifier: Callable[[AgentCard], Awaitable[AgentCard] | AgentCard]| None = None,
extended_card_modifier: Callable[[AgentCard, ServerCallContext], Awaitable[AgentCard] | AgentCard]|None = None,
max_content_length: int | None = 10 * 1024 * 1024,
) -> None
async def _handle_requests(self, request: Request) -> Response
@abstractmethod
def build(
self,
agent_card_url: str = AGENT_CARD_WELL_KNOWN_PATH,
rpc_url: str = DEFAULT_RPC_URL,
extended_agent_card_url: str = EXTENDED_AGENT_CARD_PATH,
**kwargs: Any,
) -> FastAPI | Starlette
JSONRPCApplication的构造函数定义了如下的参数:
- agent_card:
AgentCard对象,包含了Agent的身份信息和功能描述; - http_handler:
RequestHandler对象,负责处理A2A协议定义的各种请求; - extended_agent_card:可选的
AgentCard对象,包含了Agent的扩展身份信息和功能描述; - context_builder:可选的
CallContextBuilder对象,用于构建调用上下文; - card_modifier:可选的函数,用于在返回
AgentCard之前对其进行修改; - extended_card_modifier:可选的函数,用于在返回扩展
AgentCard之前对其进行修改; - max_content_length:可选的整数,表示请求内容的最大长度,默认为10MB。
JSONRPCApplication针对请求的处理实现在_handle_requests方法中。由于采用JSON-RPC协议,所以上面介绍的没A2A操作都具有一个固定的Method。一旦Method被确定,请求的意图和输入出Schema就都被确定了。_handle_requests的实现很简单,它从请求中提取Method,并将主体内容分序列化成操作对应的输入,然后将后续的请求处理分发给RequestHandler对应的方法就可以了。
上面演示实例使用A2AStarletteApplication就继承自JSONRPCApplication这个抽象类,它实现的build方法会返回一个Starlette对象。JSONRPCApplication还具有如下这个名为A2AFastAPIApplication的子类,它的build方法则会返回一个FastAPI对象。FastAPI会在Starlette的基础上提供更多的功能,例如自动生成API文档、请求参数验证等。开发者可以根据自己的需求选择使用A2AStarletteApplication还是A2AFastAPIApplication来构建ASGIApplication。
class A2AStarletteApplication(JSONRPCApplication):
def build(
self,
agent_card_url: str = AGENT_CARD_WELL_KNOWN_PATH,
rpc_url: str = DEFAULT_RPC_URL,
extended_agent_card_url: str = EXTENDED_AGENT_CARD_PATH,
**kwargs: Any,
) -> Starlette:
class A2AFastAPIApplication(JSONRPCApplication):
def build(
self,
agent_card_url: str = AGENT_CARD_WELL_KNOWN_PATH,
rpc_url: str = DEFAULT_RPC_URL,
extended_agent_card_url: str = EXTENDED_AGENT_CARD_PATH,
**kwargs: Any,
) -> FastAPI:
2. 构建CallContext
ASGIApplication面向的请求上下文体现为starllet.requests.Request对象,为了方便后续处理,JSONRPCApplication会利用__init__方法提供的CallContextBuilder将其转换成一个ServerCallContext对象作为A2A服务端调用上下文。
class ServerCallContext(BaseModel):
state: State
user: User
requested_extensions: set[str]
activated_extensions: set[str]
class CallContextBuilder(ABC):
@abstractmethod
def build(self, request: Request) -> ServerCallContext
State = collections.abc.MutableMapping[str, typing.Any]
ServerCallContext提供了如下的上下文信息:
- state:一个可变的字典对象,用于在请求处理过程中存储和共享状态信息;
- user:一个
User对象,包含了请求发起者的身份信息; - requested_extensions:一个字符串集合,表示请求中包含的A2A扩展功能;
- activated_extensions:一个字符串集合,表示在请求处理过程中被激活的扩展功能。
如果在构造JSONRPCApplication对象时没有提供CallContextBuilder,那么JSONRPCApplication会使用一个默认的DefaultCallContextBuilder来构建ServerCallContext对象。DefaultCallContextBuilder会从请求中提取用户信息、请求头信息以及请求中包含的A2A扩展功能,并将它们封装成一个ServerCallContext对象返回。
class DefaultCallContextBuilder(CallContextBuilder):
"""A default implementation of CallContextBuilder."""
def build(self, request: Request) -> ServerCallContext:
user: A2AUser = UnauthenticatedUser()
state = {}
with contextlib.suppress(Exception):
user = StarletteUserProxy(request.user)
state['auth'] = request.auth
state['headers'] = dict(request.headers)
return ServerCallContext(
user=user,
state=state,
requested_extensions=get_requested_extensions(
request.headers.getlist(HTTP_EXTENSION_HEADER)
)
)
HTTP_EXTENSION_HEADER = 'X-A2A-Extensions'
3. 请求处理器RequestHandler
我们在前面系统介绍了A2A协议涉及的所有操作,包括核心的Agent调用(请求-响应形式和流式)、任务管理、推送通知配置管理和读取AgentCard(包括读取扩展AgentCard)。这些操作具有各自的输入输出Schema,并关联着一个固定的JSON-RPC Method。这些操作体现了A2A Server需要处理的请求类型,所以作为请求处理器的RequestHandler定义了对应的抽象方法来处理这些请求。
class RequestHandler(ABC):
@abstractmethod
async def on_get_task(
self,
params: GetTaskRequest,
context: ServerCallContext,
) -> Task | None
@abstractmethod
async def on_list_tasks(
self, params: ListTasksRequest, context: ServerCallContext
) -> ListTasksResponse
@abstractmethod
async def on_cancel_task(
self,
params: CancelTaskRequest,
context: ServerCallContext,
) -> Task | None
@abstractmethod
async def on_message_send(
self,
params: SendMessageRequest,
context: ServerCallContext,
) -> Task | Message
@abstractmethod
async def on_message_send_stream(
self,
params: SendMessageRequest,
context: ServerCallContext,
) -> AsyncGenerator[Event]
@abstractmethod
async def on_create_task_push_notification_config(
self,
params: TaskPushNotificationConfig,
context: ServerCallContext,
) -> TaskPushNotificationConfig
@abstractmethod
async def on_get_task_push_notification_config(
self,
params: GetTaskPushNotificationConfigRequest,
context: ServerCallContext,
) -> TaskPushNotificationConfig
@abstractmethod
async def on_subscribe_to_task(
self,
params: SubscribeToTaskRequest,
context: ServerCallContext,
) -> AsyncGenerator[Event]
@abstractmethod
async def on_list_task_push_notification_configs(
self,
params: ListTaskPushNotificationConfigsRequest,
context: ServerCallContext,
) -> ListTaskPushNotificationConfigsResponse
@abstractmethod
async def on_delete_task_push_notification_config(
self,
params: DeleteTaskPushNotificationConfigRequest,
context: ServerCallContext,
) -> None
@abstractmethod
async def on_get_extended_agent_card(
self,
params: GetExtendedAgentCardRequest,
context: ServerCallContext,
) -> AgentCard
JSONRPCApplication从请求中提取Method后,会将请求反序列化成对应操作的输入类型,并调用RequestHandler中对应的方法来处理请求。RequestHandler的具体实现由开发者提供,A2A SDK提供了一个默认的实现DefaultRequestHandler,开发者可以直接使用它或者在它的基础上进行定制化开发。
class DefaultRequestHandler(RequestHandler):
def __init__(
self,
agent_executor: AgentExecutor,
task_store: TaskStore,
queue_manager: QueueManager | None = None,
push_config_store: PushNotificationConfigStore | None = None,
push_sender: PushNotificationSender | None = None,
request_context_builder: RequestContextBuilder | None = None,
) -> None:
当我们调用__init__方法创建DefaultRequestHandler对象时,需要提供如下的参数:
- agent_executor:
AgentExecutor对象,负责执行与Agent交互的核心逻辑; - task_store:
TaskStore对象,用于存储和管理任务; - queue_manager:可选的
QueueManager对象,用于管理事件队列; - push_config_store:可选的
PushNotificationConfigStore对象,用于存储和管理推送通知的配置; - push_sender:可选的
PushNotificationSender对象,用于发送推送通知; - request_context_builder:可选的
RequestContextBuilder对象,用于构建请求上下文。
4. Agent执行器AgentExecutor
如果说Uvicorn这个Web Server是整个A2A SDK服务端框架的起点,那么AgentExecutor就是整个框架中的终点,因为我们构建这个Web Server的最终目的是为了能够远程调用Agent来完成任务。由于Agent开发平台的多样性和复杂性,Agent的执行方法各有不同,所以A2A SDK并没有直接在框架中实现与特定Agent开发平台绑定的AgentExecutor,而是提供了一个抽象的AgentExecutor类,开发者需要根据自己使用的Agent开发平台来实现这个抽象类。
AgentExecutor这个抽象类是A2A SDK的核心执行接口。它利用抽象方法execute定义了一个Agent如何接收指令、处理任务以及汇报进度的标准行为。可以把它理解为Agent的处理器(CPU),它不负责网络传输,只负责纯粹的业务逻辑。核心设计逻辑这个接口采用的是 异步+事件驱动 的模式。Agent并不直接返回结果,而是通过event_queue(事件队列)持续向外通知它的状态。cancel方法则提供了一个接口用于取消正在执行的任务。
class AgentExecutor(ABC):
@abstractmethod
async def execute(
self, context: RequestContext, event_queue: EventQueue
) -> None
@abstractmethod
async def cancel(
self, context: RequestContext, event_queue: EventQueue
) -> None
EventQueue是A2ASDK中的通信枢纽和缓冲器。它位于执行任务的Agent(AgentExecutor)和接收结果的服务器之间,专门负责异步事件的流式传输。它采用流式传输的蓄水池设计来解决Agent执行任务通常很慢,但用户希望看到实时反馈之间的矛盾。EventQueue允许Agent边做边发(enqueue_event),服务端边收边转(dequeue_event),从而实现类似SSE的流式响应。它通过max_queue_size限制队列大小。如果队列满了,enqueue_event操作会挂起,这能防止Agent产生事件的速度远超消费速度,导致内存溢出。EventQueue通过QueueManager进行管理,后者在构建DefaultRequestHandler时指定。
class EventQueue:
def __init__(self, max_queue_size: int = DEFAULT_MAX_QUEUE_SIZE) -> None
async def enqueue_event(self, event: Event) -> None
async def dequeue_event(self, no_wait: bool = False) -> Event
def task_done(self) -> None
def tap(self) -> EventQueue
async def close(self, immediate: bool = False) -> None
def is_closed(self) -> bool
async def clear_events(self, clear_child_queues: bool = True) -> None
Event = Message | Task | TaskStatusUpdateEvent | TaskArtifactUpdateEvent
添加到EventQueue所谓的事件(Event)可以是Message、Task、TaskStatusUpdateEvent或者TaskArtifactUpdateEvent中的任意一种。还记得以流式执行Agent的SendStreamMessage操作吗?它实时响应的类型StreamReponse的字段payload所属的类型就是这四个类型之一。
async def on_message_send_stream(
self,
params: SendMessageRequest,
context: ServerCallContext,
) -> AsyncGenerator[Event]
execute和cancel方法从作为context参数的RequestContext对象中获取当前请求的上下文信息,例如当前任务、相关任务、请求消息等。RequestContext是A2A SDK中一个非常重要的类,它封装了与当前请求相关的所有上下文信息,并提供了一系列的方法和属性来访问这些信息。AgentExecutor通过RequestContext来获取输入、更新状态以及访问与当前请求相关的各种信息。RequestContext由RequestContextBuilder构建,后者在构建DefaultRequestHandler时指定。
class RequestContext:
def __init__( # noqa: PLR0913
self,
call_context: ServerCallContext,
request: SendMessageRequest | None = None,
task_id: str | None = None,
context_id: str | None = None,
task: Task | None = None,
related_tasks: list[Task] | None = None,
task_id_generator: IDGenerator | None = None,
context_id_generator: IDGenerator | None = None,
)
def get_user_input(self, delimiter: str = '\n') -> str
def attach_related_task(self, task: Task) -> None
@property
def message(self) -> Message | None
@property
def related_tasks(self) -> list[Task]
@property
def current_task(self) -> Task | None
@current_task.setter
def current_task(self, task: Task | None) -> None
@property
def task_id(self) -> str | None
@property
def context_id(self) -> str | None
@property
def configuration(self) -> SendMessageConfiguration | None
@property
def call_context(self) -> ServerCallContext
@property
def metadata(self) -> dict[str, Any]
@property
def tenant(self) -> str
@property
def requested_extensions(self) -> set[str]
RequestContext利用属性提供如下的上下文信息:
- message:一个
Message对象,表示当前请求的消息内容; - related_tasks:一个
Task对象列表,表示与当前请求相关的任务; - current_task:一个
Task对象,表示当前正在处理的任务; - task_id:一个字符串,表示当前任务的唯一标识符;
- context_id:一个字符串,表示当前请求上下文的唯一标识符;
- configuration:一个
SendMessageConfiguration对象,表示当前请求的配置信息; - call_context:一个
ServerCallContext对象,表示当前请求的调用上下文; - metadata:一个字典,表示当前请求的元数据信息;
- tenant:一个字符串,表示当前请求所属的租户信息;
- requested_extensions:一个字符串集合,表示当前请求中包含的A2A扩展功能;
RequestContext还提供了get_user_input方法来获取用户输入的内容,以及attach_related_task方法来将一个任务与当前请求相关联。AgentExecutor可以通过这些属性和方法来访问与当前请求相关的各种信息,并根据需要进行处理。
5. 其他辅助对象
在构建作为默认请求处理器的DefaultRequestHandler时,除了指定用来执行Agent的AgentExecutor之外,我们还可以指定:
- TaskStore:用于存储和管理任务的对象;
- QueueManager:用于管理事件队列的对象;
- PushNotificationConfigStore:用于存储和管理推送通知配置的对象;
- PushNotificationSender:用于发送推送通知的对象;
- RequestContextBuilder:用于构建请求上下文的对象;
5.1 TaskStore
TaskStore是一个抽象类,定义了用于存储和管理任务的接口。它提供了save、get、list和delete等方法来保存、获取、列出和删除任务。开发者需要根据自己的需求来实现这个抽象类,例如可以使用内存存储、数据库存储或者分布式存储等方式来实现TaskStore。系统提供了一个InMemoryTaskStore的实现,它使用一个字典来存储任务,适用于简单的场景和测试环境。
class TaskStore(ABC):
@abstractmethod
async def save(self, task: Task, context: ServerCallContext) -> None
@abstractmethod
async def get(
self, task_id: str, context: ServerCallContext
) -> Task | None
@abstractmethod
async def list(
self,
params: ListTasksRequest,
context: ServerCallContext,
) -> ListTasksResponse
@abstractmethod
async def delete(self, task_id: str, context: ServerCallContext) -> None
class InMemoryTaskStore(TaskStore)
5.2 QueueManager
AgentExecutor通过EventQueue来向外发送事件,而EventQueue的管理则是由QueueManager来负责的。QueueManager是一个抽象类,定义了用于管理事件队列的接口。它提供了add、get、tap、close和create_or_tap等方法来添加、获取、复制、关闭和创建事件队列。开发者需要根据自己的需求来实现这个抽象类,例如可以使用内存存储或者分布式存储等方式来实现QueueManager。系统提供了一个InMemoryQueueManager的实现,它使用一个字典来存储事件队列,适用于简单的场景和测试环境。
class QueueManager(ABC):
@abstractmethod
async def add(self, task_id: str, queue: EventQueue) -> None
@abstractmethod
async def get(self, task_id: str) -> EventQueue | None
@abstractmethod
async def tap(self, task_id: str) -> EventQueue | None
@abstractmethod
async def close(self, task_id: str) -> None
@abstractmethod
async def create_or_tap(self, task_id: str) -> EventQueue
class InMemoryQueueManager(QueueManager):
5.3 PushNotificationConfigStore
对于执行Agent这种长耗时操作,维护长连接来实现流式响应是一个常见的做法,但有些场景下可能无法使用长连接,例如网络环境不稳定或者客户端不支持长连接等。A2A采用webhook的形式来提供推送通知的功能,以解决长连接无法使用的场景。A2A SDK通过PushNotificationConfigStore这个抽象类来定义推送通知配置的存储和管理接口。它提供了set_info、get_info和delete_info等方法来设置、获取和删除推送通知配置。开发者需要根据自己的需求来实现这个抽象类,例如可以使用内存存储或者数据库存储等方式来实现PushNotificationConfigStore。系统提供了一个InMemoryPushNotificationConfigStore的实现,它使用一个字典来存储推送通知配置,适用于简单的场景和测试环境。
class PushNotificationConfigStore(ABC):
@abstractmethod
async def set_info(self, task_id: str, notification_config: PushNotificationConfig) -> None
@abstractmethod
async def get_info(self, task_id: str) -> list[PushNotificationConfig]
@abstractmethod
async def delete_info(self, task_id: str, config_id: str | None = None) -> None
class BasePushNotificationSender(PushNotificationSender):
def __init__(
self,
httpx_client: httpx.AsyncClient,
config_store: PushNotificationConfigStore,
) -> None:
self._client = httpx_client
self._config_store = config_store
async def send_notification(self, task: Task) -> None
class InMemoryPushNotificationConfigStore(PushNotificationConfigStore):
5.4 PushNotificationSender
服务端向客户端发送推送通知的功能由PushNotificationSender这个抽象类定义。它提供了一个send_notification方法,接受一个任务ID和一个事件作为参数,用于向客户端发送推送通知。开发者需要根据自己的需求来实现这个抽象类,例如可以使用HTTP请求或者消息队列等方式来实现PushNotificationSender。系统提供了一个BasePushNotificationSender的实现,它使用 httpx.AsyncClient来发送HTTP请求,并通过PushNotificationConfigStore来获取推送通知的配置。
class PushNotificationSender(ABC):
@abstractmethod
async def send_notification(self, task_id: str, event: PushNotificationEvent) -> None
PushNotificationEvent = Task | TaskStatusUpdateEvent | TaskArtifactUpdateEvent
class BasePushNotificationSender(PushNotificationSender):
def __init__(
self,
httpx_client: httpx.AsyncClient,
config_store: PushNotificationConfigStore,
) -> None
async def send_notification(self, task: Task) -> None
5.5 RequestContextBuilder
AgentExecutor的execute和cancel方法接受一个RequestContext对象作为参数,RequestContext封装了与当前请求相关的所有上下文信息。RequestContext由RequestContextBuilder构建,RequestContextBuilder是一个抽象类,定义了用于构建请求上下文的接口。它提供了一个build方法,接受一个ServerCallContext对象和一些可选的参数,用于构建一个RequestContext对象。开发者需要根据自己的需求来实现这个抽象类,例如可以从ServerCallContext中提取用户输入、相关任务等信息来构建RequestContext。系统提供了一个SimpleRequestContextBuilder的实现,它可以选择是否从TaskStore中获取与当前请求相关的任务,并将它们包含在构建的RequestContext中。
class RequestContextBuilder(ABC):
@abstractmethod
async def build(
self,
context: ServerCallContext,
params: SendMessageRequest | None = None,
task_id: str | None = None,
context_id: str | None = None,
task: Task | None = None,
) -> RequestContext
class SimpleRequestContextBuilder(RequestContextBuilder):
def __init__(
self,
should_populate_referred_tasks: bool = False,
task_store: TaskStore | None = None,
task_id_generator: IDGenerator | None = None,
context_id_generator: IDGenerator | None = None,
) -> None:
self._task_store = task_store
self._should_populate_referred_tasks = should_populate_referred_tasks
self._task_id_generator = task_id_generator
self._context_id_generator = context_id_generator
async def build(
self,
params: MessageSendParams | None = None,
task_id: str | None = None,
context_id: str | None = None,
task: Task | None = None,
context: ServerCallContext | None = None,
) -> RequestContext:
related_tasks: list[Task] | None = None
if (
self._task_store
and self._should_populate_referred_tasks
and params
and params.message.reference_task_ids
):
tasks = await asyncio.gather(
*[
self._task_store.get(task_id)
for task_id in params.message.reference_task_ids
]
)
related_tasks = [x for x in tasks if x is not None]
return RequestContext(
request=params,
task_id=task_id,
context_id=context_id,
task=task,
related_tasks=related_tasks,
call_context=context,
task_id_generator=self._task_id_generator,
context_id_generator=self._context_id_generator,
)
openEuler 是由开放原子开源基金会孵化的全场景开源操作系统项目,面向数字基础设施四大核心场景(服务器、云计算、边缘计算、嵌入式),全面支持 ARM、x86、RISC-V、loongArch、PowerPC、SW-64 等多样性计算架构
更多推荐

所有评论(0)