自托管 A2A 代理

使用 .NET A2A 托管包通过 ASP.NET Core 公开 Agent Framework 代理。 有关包设置和完整的服务器示例,请参阅 A2A 集成

将 Go provider/a2aprovider 包与官方 A2A Go 服务器处理程序配合使用。 有关完整的服务器示例,请参阅 A2A 集成

代理框架通过官方 A2A SDK 提供两个用于托管代理和工作流的Python包:

Package 集成模型 在以下情况下使用
agent-framework-a2a 一个具有强约定的 A2AExecutor,可转换请求、运行代理,并发布 A2A 任务事件和工件。 您希望采用标准的 Agent Framework 到 A2A 的行为,并且只需要搭建 A2A SDK 服务器。
agent-framework-hosting-a2a 用于应用自有执行器的增量式构建模块。 从基础的代理或工作流转换器开始,并可选择使用构建在这些转换器之上的 AgentA2AAdapterWorkflowA2AAdapter,以添加原生卡片生成功能和模式验证。 应用程序需要拥有会话映射、任务转换、事件传递、项目边界、输出转换或多协议主机。

这两个包都使用原生 A2A SDK 的类型和服务器组件。 应用程序提供请求处理程序、任务存储、路由或 SDK 应用程序生成器、身份验证和部署。 应用程序 agent-framework-hosting-a2a可以直接构造代理卡,也可以让适配器生成它。

使用有意见的 A2A 执行程序

当内置服务器适配器符合你的生命周期需求时,安装 agent-framework-a2a

pip install --pre agent-framework-a2a starlette uvicorn

A2AExecutor 实现 A2A SDK 的 AgentExecutor。 它从 A2A 请求上下文读取用户输入,从 A2A 上下文 ID 创建代理框架会话,在流式处理或非流式处理模式下运行代理,转换支持的输出内容,并通过 SDK TaskUpdater发布任务状态和项目事件。

使用 A2A SDK、 DefaultRequestHandler任务存储、代理卡和 Starlette 应用程序或其他受支持的服务器集成编写它。 使用 A2AExecutor(agent, stream=True) 配置流式处理,通过 run_kwargs 传递稳定的代理运行选项,或者在需要不同的输出映射时,继承 A2AExecutor 并重写 handle_events

A2AExecutor 范围限定为 A2A 终结点,并直接管理其 A2A 执行和会话映射。 当同一代理必须通过一个应用程序中的多个协议可用时,请使用托管包。

有关完整的服务器设置,请参阅 通过 A2A 公开代理框架代理

在应用拥有的执行程序中使用适配器

当应用程序拥有原生 A2A 执行器,但希望 Agent Framework 生成公开卡片并验证转换时,请安装托管包:

pip install --pre agent-framework-hosting-a2a starlette uvicorn

AgentA2AAdapter 支持代理或 AgentState。 其异步 get_card 方法会获取公开名称和描述,默认采用较为保守的文本模式,并且可以从 Agent Framework SkillsProvider 实例中推断出原生 A2A 技能。 服务器功能和支持的接口保持显式,因为它们描述应用程序终结点,而不是代理 run 的方法。

适配器公开 a2a_to_runa2a_from_run 方法,这些方法默认针对配置的卡模式验证值。 应用程序仍拥有 A2A 执行程序、任务生命周期、事件队列、项目边界、会话策略、身份验证、路由和部署。

此执行程序使用一个适配器进行入站转换、代理状态和出站转换:

class AppAgentExecutor(AgentExecutor):
    """Native A2A SDK executor composed with Agent Framework conversion helpers."""

    def __init__(self, adapter: AgentA2AAdapter[Any]) -> None:
        self.adapter = adapter

    async def cancel(self, context: RequestContext, event_queue: EventQueue) -> None:
        if context.context_id is None:
            raise ValueError("A2A context id is required")
        updater = TaskUpdater(event_queue, context.task_id or "", context.context_id)
        await updater.cancel()

    async def execute(self, context: RequestContext, event_queue: EventQueue) -> None:
        if context.message is None or context.context_id is None:
            raise ValueError("A2A message and context id are required")

        task = context.current_task
        if task is None:
            task = new_task_from_user_message(context.message)
            await event_queue.enqueue_event(task)

        updater = TaskUpdater(event_queue, task.id, context.context_id)
        await updater.submit()
        try:
            await updater.start_work()
            run = self.adapter.a2a_to_run(context.message, stream=True)
            agent = await self.adapter.state.get_target()
            # Demo-only key: the outer server must authenticate and authorize these protocol IDs for multi-user use.
            session_id = f"a2a:{context.tenant}:{context.context_id}"
            session = await self.adapter.state.get_or_create_session(session_id)
            if not run["stream"]:
                raise RuntimeError("This executor requires streaming run arguments.")
            stream = agent.run(  # pyright: ignore[reportCallIssue]
                run["messages"],
                session=session,
                options=run["options"],
                stream=run["stream"],
            )
            default_artifact_id = uuid.uuid4().hex
            streamed_artifact_ids: set[str] = set()
            async for update in stream:
                parts = self.adapter.a2a_from_run(update)
                if parts:
                    artifact_id = update.message_id or default_artifact_id
                    await updater.add_artifact(
                        parts=parts,
                        artifact_id=artifact_id,
                        append=True if artifact_id in streamed_artifact_ids else None,
                    )
                    streamed_artifact_ids.add(artifact_id)
            final_response = await stream.get_final_response()
            if not streamed_artifact_ids:
                parts = self.adapter.a2a_from_run(final_response)
                if parts:
                    await updater.update_status(
                        state=TaskState.TASK_STATE_WORKING,
                        message=updater.new_agent_message(parts),
                    )
            await self.adapter.state.set_session(session_id, session)
            await updater.complete()
        except asyncio.CancelledError:
            await updater.update_status(state=TaskState.TASK_STATE_CANCELED)
        except Exception:
            logger.exception("A2A agent execution failed.")
            await updater.update_status(
                state=TaskState.TASK_STATE_FAILED,
                message=updater.new_agent_message([Part(text="Agent execution failed.")]),
            )

服务器设置会创建适配器,生成其原生 AgentCard,并将应用自有执行器与 A2A SDK 请求处理程序组合在一起:

if __name__ == "__main__":
    flight_skill = InlineSkill(
        frontmatter=SkillFrontmatter(
            name="flight-booking",
            description="Search and book flights across Europe.",
        ),
        instructions="Help users search and book flights across Europe.",
    )
    hotel_skill = InlineSkill(
        frontmatter=SkillFrontmatter(
            name="hotel-booking",
            description="Search and book hotels across Europe.",
        ),
        instructions="Help users search and book hotels across Europe.",
    )
    agent = Agent(
        client=OpenAIChatClient(),
        name="Europe Travel Agent",
        description="Helps users search and book flights and hotels across Europe.",
        instructions="You are a helpful Europe Travel Agent.",
        context_providers=[SkillsProvider([flight_skill, hotel_skill])],
    )

    state = AgentState(agent)
    adapter = AgentA2AAdapter(
        state,
        version="1.0.0",
        capabilities=AgentCapabilities(streaming=True),
        supported_interfaces=[AgentInterface(url="http://localhost:9999/", protocol_binding="JSONRPC")],
    )
    public_agent_card = asyncio.run(adapter.get_card())
    request_handler = DefaultRequestHandler(
        agent_executor=AppAgentExecutor(adapter),
        task_store=InMemoryTaskStore(),
        agent_card=public_agent_card,
    )

构建一个应用专属的 A2A 执行器

当应用程序还需要直接控制卡创建时,请使用独立托管帮助程序:

pip install --pre agent-framework-hosting-a2a starlette uvicorn

这些辅助工具与框架无关:

  • a2a_to_run 将 A2A Message 转换为代理框架运行参数。
  • a2a_from_run 将 Agent Framework 响应和流式更新转换为 A2A Part 值。

你的执行器会选择会话密钥,并负责管理任务状态转换、事件队列、构件 ID、消息边界和出站交付。 a2a_from_run 返回平面部件列表,以便应用程序可以将这些部件分组到 A2A 消息或项目中,并应用消息级元数据。

托管设置还支持多协议应用程序。 在 A2A、OpenAI 响应、Telegram 和 MCP 路由之间共享相同的代理目标和 AgentState 基础结构,而每个协议终结点都保留自己的转换、授权和会话密钥策略。 这样,客户端就可以同时通过不同的协议访问一个代理,而无需为每个终结点创建单独的代理部署。

在原生 A2A SDK 执行器中组合辅助函数。 此示例创建并更新 A2A 任务,将入站消息转换为 Agent Framework 运行实例,在流结束后持久保存更新后的 AgentState 会话,并将返回的部分内容作为工件发布。

class AppAgentExecutor(AgentExecutor, Generic[AgentT]):
    """Native A2A SDK executor composed with Agent Framework conversion helpers."""

    def __init__(self, state: AgentState[AgentT]) -> None:
        self.state = state

    async def cancel(self, context: RequestContext, event_queue: EventQueue) -> None:
        if context.context_id is None:
            raise ValueError("A2A context id is required")
        updater = TaskUpdater(event_queue, context.task_id or "", context.context_id)
        await updater.cancel()

    async def execute(self, context: RequestContext, event_queue: EventQueue) -> None:
        if context.message is None or context.context_id is None:
            raise ValueError("A2A message and context id are required")

        task = context.current_task
        if task is None:
            task = new_task_from_user_message(context.message)
            await event_queue.enqueue_event(task)

        updater = TaskUpdater(event_queue, task.id, context.context_id)
        await updater.submit()
        try:
            await updater.start_work()
            run = a2a_to_run(context.message, stream=True)
            agent = await self.state.get_target()
            # Demo-only key: the outer server must authenticate and authorize these protocol IDs for multi-user use.
            session_id = f"a2a:{context.tenant}:{context.context_id}"
            session = await self.state.get_or_create_session(session_id)
            if not run["stream"]:
                raise RuntimeError("This executor requires streaming run arguments.")
            stream = agent.run(  # pyright: ignore[reportCallIssue]
                run["messages"],
                session=session,
                options=run["options"],
                stream=run["stream"],
            )
            default_artifact_id = uuid.uuid4().hex
            streamed_artifact_ids: set[str] = set()
            async for update in stream:
                parts = a2a_from_run(update)
                if parts:
                    artifact_id = update.message_id or default_artifact_id
                    await updater.add_artifact(
                        parts=parts,
                        artifact_id=artifact_id,
                        append=True if artifact_id in streamed_artifact_ids else None,
                    )
                    streamed_artifact_ids.add(artifact_id)
            final_response = await stream.get_final_response()
            if not streamed_artifact_ids:
                parts = a2a_from_run(final_response)
                if parts:
                    await updater.update_status(
                        state=TaskState.TASK_STATE_WORKING,
                        message=updater.new_agent_message(parts),
                    )
            await self.state.set_session(session_id, session)
            await updater.complete()
        except CancelledError:
            await updater.update_status(state=TaskState.TASK_STATE_CANCELED)
        except Exception:
            logger.exception("A2A agent execution failed.")
            await updater.update_status(
                state=TaskState.TASK_STATE_FAILED,
                message=updater.new_agent_message([Part(text="Agent execution failed.")]),
            )

该示例使用 Starlette 和 Uvicorn,但帮助程序也不与这两者绑定。 使用您的应用框架或 A2A SDK 应用构建器来提供 A2A 代理卡和 JSON-RPC 路由:

# Create the Agent Framework agent for the chosen type
agent_factory = AGENT_FACTORIES[args.agent_type]
agent = agent_factory(client)
state = AgentState(agent)

# Build the A2A server components
url = f"http://{args.host}:{args.port}/"
agent_card = AGENT_CARD_FACTORIES[args.agent_type](url)
executor = AppAgentExecutor(state)
task_store = InMemoryTaskStore()
request_handler = DefaultRequestHandler(
    agent_executor=executor,
    task_store=task_store,
    agent_card=agent_card,
)

app = Starlette(
    routes=[
        *create_agent_card_routes(agent_card),
        *create_jsonrpc_routes(request_handler, "/"),
    ]
)

使用适配器托管工作流

WorkflowA2AAdapter 为工作流或 WorkflowState 提供相同的卡片生成和转换边界。 它从工作流的声明类型推断保守的输入和输出模式,也可以为特定于应用程序的表示形式提供显式模式。

独立的 a2a_to_workflow_runa2a_from_workflow_run 辅助程序用于对工作流输入和输出进行类型化转换。 适配器将其公开为异步 a2a_to_run 和同步 a2a_from_run 方法,以针对其有效卡模式进行验证。 输入转换接受工作流单一启动执行器输入类型的一个 A2A 文本、原始或数据部件,输出转换会将已完成的公共工作流输出映射到本机 A2A 部件。 当适配器必须推断输出模式时,在验证的输出转换之前调用 get_card

应用程序仍负责原生 A2A 执行器,以及对进度、任务状态、工件、检查点和人在回路中的后续处理进行流式传输。 处于挂起状态的人工输入请求不会自动转为继续执行,因此主机必须自行实现续接策略。

保护会话和任务状态

A2AExecutor 使用 A2A 上下文 ID 作为代理框架会话 ID。 基于适配器和帮助程序的示例将 A2A 租户和上下文 ID 结合起来,用于演示由应用程序选择的映射关系。 无论采用哪种方案,生产主机都必须在调用方到达 A2A 请求处理程序之前先对其完成身份验证,并根据该受信任身份确定租户和主体,同时对所有任务、上下文、继续和取消 ID 进行授权校验。

Important

A2A SDK 的默认任务存储和推送配置存储位于内存中,并按用户名划分所有权范围。 对于多租户服务,请使用从同一受信任租户和主体派生所有权的 owner_resolver,并在副本可能重启或横向扩展时使用持久化任务存储和会话存储。

有关基于帮助程序的完整服务器和多代理示例,请参阅 A2A 托管示例。 有关 A2A 客户端和协议功能,请参阅 A2A 集成

后续步骤

更深入: