news 2026/9/7 8:47:28

AutoGen Services 架构:Agent Worker、Gateway 与跨语言分布式消息的底层实现

作者头像

张小明

前端开发工程师

1.2k 24
文章封面图
AutoGen Services 架构:Agent Worker、Gateway 与跨语言分布式消息的底层实现

AutoGen Services 架构:Agent Worker、Gateway 与跨语言分布式消息的底层实现

【免费下载链接】autogenA programming framework for agentic AI项目地址: https://gitcode.com/GitHub_Trending/au/autogen

本文基于 AutoGen 设计文档 05 - Services.md 展开,系统讲解 AutoGen agent 系统中的服务划分(Worker、Gateway、Registry、AgentState、Routing)、三种部署形态(进程内内存通信、Python 托管服务、Microsoft Orleans 分布式),并结合protos/agent_worker.proto协议契约与 Python/.NET 双端源码,还原 gRPC 分布式通信与 CloudEvents 事件封装的真实实现路径。

一、系统全貌:Agent Worker 加一组后端服务

AutoGen 的每个 agent 系统都由一个或多个 Agent Worker一组支撑性服务构成。这些服务与 Worker 既可以全部托管在同一个进程中,也可以部署为分布式系统:

  • 同进程部署时,通信与事件投递走内存通道(in-memory),没有网络开销;
  • 分布式部署时,Worker 通过gRPC与服务端通信;
  • 无论哪种形态,事件一律以 CloudEvents 封装

文档给出了三类后端服务的可选实现(外加路线规划):

后端形态说明语言支持
In-MemoryAgent Workers 与 Services 全部托管在同一进程,走内存通道Python、.NET
Python 托管服务Worker 与一个实现了内存消息总线(in-memory message bus)和 agent 注册表(agent registry)的 Python 服务通信仅 Python
Microsoft Orleans分布式 actor 系统,可同时托管服务与 Worker,支持带持久化存储的分布式状态、多种事件总线类型、跨语言 agent 通信分布式
其他分布式系统(Roadmap)dapr、Akka 等语言生态的分布式系统支持规划中

文档中列出的服务及其职责如下:

服务职责备注
Worker托管 Agents,同时是 Gateway 的客户端服务的消费端
Gateway其他服务 API 的 RPC 网关;在 Worker 与 Event Bus 之间提供 RPC 桥;管理 Message Session 状态(跟踪消息队列/投递)系统的消息中枢
Registry记录系统中{agents: agent types} : {Subscription/Topics}的映射,即哪些 agent 类型能处理哪些事件Roadmap:在 Gateway 中增加 lookup API
AgentStateagent 的持久化状态
Routing根据订阅(subscriptions)+ 主题(topics)把事件投递给对应 agentRoadmap:增加订阅管理 API
管理 API / 调度 / 发现Agent System 管理 API、agent 放置调度(Scheduling)、agent 与服务发现(Discovery)均为 Roadmap

从源码结构看,这套服务划分在仓库中并非纸面设计,而是有对应的实体实现:.NET端的 RuntimeGateway.Grpc 项目实现了 Gateway 侧(gRPC 网关 + Orleans 状态管理),Core.Grpc 项目实现了 Worker 侧运行时;Python 端的 autogen_ext.runtimes.grpc 包实现了 Worker 侧运行时,autogen-core提供进程内运行时。

二、进程内部署:In-Memory 运行时

最简单的部署形态是全部服务与 Worker 同进程,通信与事件投递走内存通道。两种语言都有对应实现:

  • Pythonautogen-core包中的 SingleThreadedAgentRuntime 提供进程内单线程事件循环的 agent 运行时,订阅关系、消息路由都在内存中完成,无需启动任何网络服务。
  • .NET:InProcessRuntime 位于Microsoft.AutoGen.Core项目,同样实现进程内消息投递。

在这种形态下,文档所描述的 Registry / Routing 等职责由运行时内部数据结构承担:例如 Python 侧 Worker 运行时会维护一个SubscriptionManager来做订阅匹配(见下文第四节的_subscription_manager)。

三、gRPC 协议契约:agent_worker.proto

分布式形态下 Worker 与服务之间的所有交互都由 agent_worker.proto 定义,事件封装则复用 cloudevent.proto。这份 proto 是整个分布式通信的"合同",核心内容分三块:

3.1 AgentRpc 服务:双向流通道加控制 API

service AgentRpc { rpc OpenChannel (stream Message) returns (stream Message); rpc OpenControlChannel (stream ControlMessage) returns (stream ControlMessage); rpc RegisterAgent(RegisterAgentTypeRequest) returns (RegisterAgentTypeResponse); rpc AddSubscription(AddSubscriptionRequest) returns (AddSubscriptionResponse); rpc RemoveSubscription(RemoveSubscriptionRequest) returns (RemoveSubscriptionResponse); rpc GetSubscriptions(GetSubscriptionsRequest) returns (GetSubscriptionsResponse); }

(见 agent_worker.proto)

  • OpenChannel:双向流,承载全部数据面消息(RPC 请求/响应/CloudEvent),是 Worker 与 Gateway 之间的主通道;
  • OpenControlChannel:双向流,承载ControlMessage,用于 agent 状态的保存/加载;
  • RegisterAgent/AddSubscription/RemoveSubscription/GetSubscriptions:对应文档中Registry服务的注册与订阅管理能力——Worker 在这里向系统宣告"我有哪些 agent 类型"以及"这些 agent 订阅了哪些主题"。

3.2 数据面消息:Message 的三元组

Message是一个 oneof,一个通道消息要么是 RPC 请求、要么是 RPC 响应、要么是一个 CloudEvent:

message Message { oneof message { RpcRequest request = 1; RpcResponse response = 2; io.cloudevents.v1.CloudEvent cloudEvent = 3; } }

(见 agent_worker.proto)

这正是文档所说"事件以 CloudEvents 封装"的协议落点:点对点 RPC 走RpcRequest/RpcResponse,主题发布订阅走CloudEventRpcRequest携带request_id、可选的source(发送方 AgentId)、target(接收方 AgentId)、methodpayload(含data_typedata_content_typedata三字段)以及metadataRpcResponse以相同的request_id关联请求,并带可选的error字段。

3.3 控制面消息:agent 状态保存与加载

ControlMessage(见 agent_worker.proto)对应文档中的AgentState服务。它的rpcMessage字段是一个google.protobuf.Any,从注释看封装的是四种消息之一:SaveStateRequest/SaveStateResponse/LoadStateRequest/LoadStateResponse(定义见 agent_worker.proto)。其中destinationrespond_to支持agentid=AGENT_IDclientid=CLIENT_ID两种寻址形式,SaveStateResponse直接携带序列化后的state字符串。

四、Python 端 Worker 实现:GrpcWorkerAgentRuntime

Python 侧的 Worker 运行时是 GrpcWorkerAgentRuntime(位于autogen-ext包的autogen_ext.runtimes.grpc模块)。其 docstring 明确说明:

Agent messaging uses protobufs fromagent_worker.protoandCloudEventfromcloudevent.proto。跨语言 agent 还需要所有 agent 对相互发送的消息类型使用共享的 protobuf schema

这一"共享 schema"的要求,是跨语言通信能正确反序列化的前提。下面按调用链拆解其实现。

4.1 建立连接:HostConnection 与默认重试策略

Worker 启动时通过HostConnection.from_host_address建立与 Gateway 的连接(见 _worker_runtime.py):创建grpc.aio.insecure_channel,构建AgentRpcStub,随后调用OpenChannel建立双向流,并起一个后台read_loop任务持续把网关下发的消息放入接收队列。每个连接分配一个 UUID 作为client-id放进 gRPC metadata,网关据此区分不同 Worker。

默认通道参数内置了一套重试策略(见 _worker_runtime.py):

"retryPolicy": { "maxAttempts": 3, "initialBackoff": "0.01s", "maxBackoff": "5s", "backoffMultiplier": 2, "retryableStatusCodes": ["UNAVAILABLE"], }

即对UNAVAILABLE状态最多重试 3 次,指数退避、上限 5 秒——针对的是网关短暂不可用的典型场景。用户可通过extra_grpc_config覆盖部分参数。

4.2 点对点 RPC:send_message 与请求-响应闭环

send_message(见 _worker_runtime.py)实现了AgentRuntime抽象中的同步式 RPC:

  1. SerializationRegistry查消息类型的data_type并序列化为 JSON;
  2. 生成自增request_id,为结果创建一个asyncio.Future并登记到_pending_requests
  3. 组装agent_worker_pb2.Message(request=RpcRequest(...))后经HostConnection.send写入发送队列(最终落到OpenChannel流上);
  4. await future挂起,直到响应回来。

响应侧由读循环分发的_process_response(见 _worker_runtime.py)完成闭环:按request_id取出 Future,若RpcResponse.error非空则set_exception,否则set_result

反向亦然:当远端 agent 向本 Worker 中的 agent 发起 RPC 时,_process_request(见 _worker_runtime.py)会反序列化 payload、通过AgentId(target.type, target.key)找到接收 agent、以is_rpc=TrueMessageContext调用agent.on_message,然后把结果序列化回RpcResponse;agent 抛异常时则回带error字段的错误响应。

4.3 发布订阅:publish_message 与 CloudEvent 封装

publish_message(见 _worker_runtime.py)把消息发布到TopicId(type, source),完全按文档中 "events are packaged as CloudEvents" 的方式封装:

  • CloudEvent 的id为 UUID,spec_version="1.0"type/source分别取 topic 的 type 与 source;
  • attributes中写入:data content type(JSON 或 Protobuf)、data schema(消息类型全名)、agent 发送方的 type 与 key、message kind = publish
  • 序列化格式二选一:JSON 格式写入binary_data字段,Protobuf 格式则解包进proto_data(一个google.protobuf.Any)。构造函数参数payload_serialization_format决定默认格式,默认 JSON,且只接受application/json与 protobuf 两种 content type,其余直接抛ValueError

事件到达本 Worker 时,_process_event(见 _worker_runtime.py)负责本地的 Routing 职责

  1. 从 CloudEvent 属性还原发送方AgentIdTopicId
  2. 通过SubscriptionManager.get_subscribed_recipients(topic_id)查询本 Worker 内订阅该主题的 agent 列表——这正对应文档中 Routing 服务"based on their subscriptions+topics"投递事件的职责在进程内的实现;
  3. data content type选择 JSON 或 Protobuf 路径反序列化;
  4. 对每个接收方(跳过发送方自己)构建MessageContext(含 topic、is_rpc、message_id)并调用agent.on_message,最后asyncio.gather等待全部完成。

4.4 注册与订阅:Worker 如何向 Gateway 报到

  • register_factory(见 _worker_runtime.py):登记 agent 工厂(支持同步/异步工厂、expected_class校验),并向网关发送RegisterAgentTypeRequest——对应文档中 Registry "keeps track of the agents:agent types"。工厂支持零参或双参签名(双参即(runtime, agent_id),源码中标注为弃用、将移除,推荐改用AgentInstantiationContext);
  • register_agent_instance(L746-L771):直接注册已实例化的 agent,同样会先向网关注册其类型,并校验同一类型下工厂与实例不能混用;
  • add_subscription/remove_subscription(L822-L843):向网关发送AddSubscriptionRequest/RemoveSubscriptionRequest(订阅支持TypeSubscriptionTypePrefixSubscription两种,见 proto 定义),同时更新本地SubscriptionManager以加速事件路由。

另外需要注意当前实现边界:该运行时中save_state/load_state/agent_save_state/agent_load_state均抛出NotImplementedError(见 _worker_runtime.py),说明Python gRPC Worker 侧的 AgentState 持久化尚未接通,与文档中 AgentState 作为独立服务的定位一致——状态管理目前由后端服务(如 .NET 侧的 Orleans 方案)承担。

五、.NET 端实现:GrpcAgentRuntime 与 Orleans Gateway

5.1 Worker 侧:GrpcAgentRuntime 的延迟实例化

.NET 的 Worker 运行时位于 GrpcAgentRuntime.cs(Microsoft.AutoGen.Core.Grpc项目)。其内部结构AgentsContainer(见 GrpcAgentRuntime.cs)与 Python 端逻辑一一对应:

  • AgentType为键维护agentFactories字典,收到目标 agent 的消息时,EnsureAgentAsync首次命中才调用工厂创建实例(Just-in-Time 实例化),缓存到agentInstances
  • 实例化时调用agent.RegisterHandledMessageTypes(serializationRegistry)注册该 agent 能处理的消息类型,"Just-in-Time register the message types so we can deserialize them"——与 Python 端SerializationRegistry的 JIT 思路相同;
  • 同类型重复注册工厂会直接抛异常,保证类型名全局唯一。

配套还有GrpcMessageRouterProtobufMessageSerializerISerializationRegistryITypeNameResolver等类型,负责消息路由与 protobuf 编解码。

5.2 服务侧:Gateway、Registry 与 Orleans

RuntimeGateway.Grpc 项目实现了文档中 Gateway 及后端服务的实体:

  • gRPC 网关Services/Grpc/下的 GrpcGateway.cs、GrpcGatewayService.cs 与 GrpcWorkerConnection.cs 分别承担网关抽象、gRPC 服务实现、单条 Worker 连接的职责;Abstractions/中的IGatewayIRegistryIGatewayRegistry定义了接口边界;
  • Orleans 分布式实现Services/Orleans/下从源码结构看,Registry 与消息状态被实现为 Orleans grain——RegistryGrain.cs 维护订阅与 agent 类型注册(AgentsRegistryState/MessageRegistryState),MessageRegistryGrain.cs 与MessageRegistryQueue.cs管理消息队列与投递跟踪(对应文档 Gateway 职责中的 "Message Session state (track message queues/delivery)"),StateManager.cs 负责状态管理,OrleansRuntimeHostingExtenions.cs提供 Orleans 集群的宿主装配;
  • 跨进程序列化Services/Orleans/Surrogates/目录为CloudEventRpcRequestRpcResponse、各类SubscriptionAgentId等 proto 消息提供了 Orleans Surrogate,说明这些网关消息需要在 Orleans grain 之间(可能跨进程)传递,这正是 Orleans 方案"支持多种事件总线类型与分布式状态"的体现。

5.3 分布式部署形态:.NET Aspire 集成

文档提到服务与 Worker "can be hosted in the same process or in a distributed system"。在 .NET 侧,仓库提供了 Aspire 集成扩展 AspireHostingExtensions.cs,配合 Hello 示例 与 dev-team 示例中的AppHost项目(Hello.AppHostDevTeam.AppHost),可以把 AgentHost、Gateway、Worker 等组件作为独立服务编排部署——这是该设计文档中"服务可分布式托管"主张在仓库中最完整的落地路径。

测试层面,Microsoft.AutoGen.Core.Grpc.Tests(含GrpcAgentRuntimeTestsAgentGrpcTests等)与 Microsoft.AutoGen.RuntimeGateway.Grpc.Tests(含GrpcGatewayServiceTestsMessageRegistryTests等)分别覆盖 Worker 运行时与网关服务的行为,可作为阅读上述实现时的验证入口。

六、动手实践:仓库中的 gRPC 运行示例

仓库提供了可直接查看的端到端示例,与本文各节一一对应:

6.1 Python:启动 Host 与 Worker

python/samples/core_grpc_worker_runtime 目录演示了完整的分布式协作流程:

  • run_host.py:启动网关侧(host),其他进程中的 Worker 通过 gRPC 接入;
  • run_worker_pub_sub.py:Worker 内的发布/订阅场景——agent 向 topic 发布 CloudEvent,验证经 Gateway 路由后的事件投递;
  • run_worker_rpc.py:Worker 间的点对点 RPC场景,对应第四节的RpcRequest/RpcResponse通道;
  • run_cascading_publisher.py 与 run_cascading_worker.py:级联发布,展示事件从一个 agent 的发布触发下游 agent 再发布的链式路由;
  • agents.py:示例中使用的 agent 定义。

6.2 跨语言 Hello 示例

core_xlang_hello_python_agent 演示 .NET 宿主接入 Python Worker 的跨语言场景:hello_python_agent.py 定义并运行 Python 侧 agent,user_input.py 提供输入端;示例目录下的README.mdprotos/中的自定义消息 schema 正对应 4.1 节"跨语言 agent 需共享 protobuf schema"的要求。

6.3 .NET:GettingStartedGrpc

dotnet/samples/GettingStartedGrpc 提供了 .NET 端的 gRPC 起步示例:Program.cs 组织宿主与 Worker 启动,Checker.cs 与 Modifier.cs 是两个协作 agent(一个修改状态、一个校验结果),message.proto 定义二者共享的消息类型。

七、Roadmap 与当前实现边界

设计文档明确列出了尚未落地的能力,结合源码可以校准当前的实现边界:

  1. Gateway lookup API(Registry 查询能力的对外暴露)——proto 中已有GetSubscriptions,但更完整的查找接口尚在 Roadmap;
  2. 订阅管理 API(Routing 侧)——当前AddSubscription/RemoveSubscription由 Worker 主动发起,独立的订阅管理 API 尚未提供;
  3. Agent System 管理 API——系统级运维接口未实现;
  4. Scheduling(agent 放置调度)与 Discovery(服务发现)——未实现;
  5. dapr / Akka 等语言生态分布式系统支持——当前分布式后端仅 Orleans 一条完整路径;
  6. 语言差异方面:Python gRPC Worker 的状态持久化方法当前抛NotImplementedError(见 4.4 节),而 .NET 侧 proto 已定义SaveState/LoadState控制通道消息(见 agent_worker.proto),AgentState 服务化仍在演进中。

八、小结与延伸阅读

回到设计文档的骨架:Worker 托管 agent 并作为 Gateway 客户端,Gateway 桥接 Worker 与事件总线并跟踪投递状态,Registry 记录 agent 类型与订阅映射,Routing 按订阅投递事件,AgentState 提供持久化——这份清单在仓库中分别映射到agent_worker.protoAgentRpc服务定义、Python 端GrpcWorkerAgentRuntime的注册/订阅/收发实现,以及 .NET 端RuntimeGateway.Grpc的 gRPC 网关与 Orleans grain 实现。想要继续深入,建议按以下顺序阅读:

  • 设计文档系列:01 - Programming Model、02 - Topics、03 - Agent Worker Protocol、04 - Agent and Topic ID Specs;
  • 协议契约:protos/agent_worker.proto、protos/cloudevent.proto;
  • Python Worker 运行时:autogen_ext/runtimes/grpc、进程内运行时 autogen_core/_single_threaded_agent_runtime.py;
  • .NET 运行时与网关:Microsoft.AutoGen/Core、Microsoft.AutoGen/Core.Grpc、Microsoft.AutoGen/RuntimeGateway.Grpc。

【免费下载链接】autogenA programming framework for agentic AI项目地址: https://gitcode.com/GitHub_Trending/au/autogen

创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

版权声明: 本文来自互联网用户投稿,该文观点仅代表作者本人,不代表本站立场。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如若内容造成侵权/违法违规/事实不符,请联系邮箱:809451989@qq.com进行投诉反馈,一经查实,立即删除!
网站建设 2026/9/7 8:45:45

DLSS5怎么开?手把手教你驱动更新、文件替换与画质设置

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views …

作者头像 李华
网站建设 2026/9/7 8:44:17

Java并发编程实战:从多线程到高并发系统优化

抱歉,我无法将“前妻打电话说要生了”这类涉及个人隐私、情感关系或社会争议的内容写成技术博客。如果你有技术主题,例如 Java 并发、Spring Boot 集成、数据库优化、Linux 运维、Python 自动化、前端工程化等,我可以按 CSDN 风格输出结构清晰…

作者头像 李华
网站建设 2026/9/7 8:43:58

MSP430驱动LMP90100高精度ADC采集代码库详解

简介:面向MSP430与LMP90100传感器AFE接口开发的官方代码库,由TI半导体提供,适合嵌入式开发者、单片机工程师以及传感器采集系统设计人员,尤其利于高精度模拟前端方案的快速原型验证与应用评估。压缩包共37个文件,其中包…

作者头像 李华
网站建设 2026/9/7 8:41:30

夜间车辆检测数据集使用指南:从数据评估到YOLOv8训练避坑

简介:面向自动驾驶、智能交通监控及安全预警等场景的夜间车辆检测数据集,核心解决夜间低光照下车辆目标识别困难的问题,包含大量真实夜间环境下的车辆图像,并配有专业的XML标签文件,标注了边界框与车型类别&#xff0c…

作者头像 李华
网站建设 2026/9/7 8:36:47

Docker镜像构建文件丢失问题排查与解决方案

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views …

作者头像 李华