跳到主要内容

Environment 消息总线、Message 路由与 Team

30 秒导读: 单个角色的「观察-思考-行动」循环(01 章)只解决了「一个 agent 怎么自转」。多个 agent 要协作,就缺一样东西:它们之间怎么通信、消息怎么被正确投递给该收的人。这一章讲 MetaGPT 的答案——Environment 是一根消息总线:角色把消息「发布」到环境,环境按一本地址簿把它投递到目标角色的私有信箱;Team 再在总线之上套一层预算 + 轮次,把整台机器驱动起来。


1. 这一章解决什么问题(零基础也能懂)

前两章的主角是单个角色:Role 是自转的原子(01 章),Action 是它手里的工作单元(02 章)。

但「多智能体」的关键从来不是「有很多角色」,而是角色之间怎么说话。设想产品经理写完了 PRD,他怎么把这份 PRD「递」给架构师、而不是递给测试工程师?架构师又怎么知道「这条消息是冲我来的、我该动手了」?

这就是本章要解决的两个问题:

  • 投递问题: 一条消息发出去,谁该收到?怎么保证精确送达、不误投、不漏投?
  • 驱动问题: 谁来推动这群角色一轮一轮地转下去,什么时候该停(活干完了 / 钱花光了)?

MetaGPT 的设计有一个关键的解耦思想(来自 RFC 113/116,见 §5):

消息只负责说清「发给谁」,完全不关心「对方在哪、怎么送过去」。 「怎么送达」是传输框架(即 Environment)的职责。

这正是 publish_message 源码注释里写的原话(metagpt/environment/base_env.py:175-183):路由信息只指定收件人,"without concern for where the message recipient is located"。

一句话直觉:Environment 想成一个公司内部的邮件系统。 你写邮件只填收件人名字(send_to),不用管对方坐哪、用什么网线——邮件系统(环境)拿着一本通讯录(member_addrs)负责把信塞进对方的收件箱(每个角色的私有 msg_buffer)。


2. 顶层全景(它大概怎么转)

三个主角

部件白话职责在哪个文件
Message一封信。带路由三元组 cause_by / sent_from / send_to,说明「因何而发、谁发的、发给谁」metagpt/schema.py:232
Environment邮件系统 / 消息总线。持有地址簿 member_addrs,负责分发并发驱动所有角色metagpt/environment/base_env.py:124
Team公司老板。雇人(hire)、给预算(invest)、按轮次运转(run),并能存档/恢复整个项目metagpt/team.py:32

一条消息怎么从发出到送达

下图是本章的主线。怎么读:从上到下是一条消息的旅程,左边是「谁在做」,右边是「用哪个符号」。

Team.run_project(idea) team.py:102 run_project
│ 发布一封「用户需求」信 Message(content=idea)

Environment.publish_message(msg) base_env.py:175

│ 遍历地址簿 member_addrs: {角色 → 它的地址集合}
│ 对每个角色问一句: is_send_to(msg, 它的地址)? common.py:423
│ ├─ msg.send_to 含 "<all>" → 命中(广播)
│ └─ 角色地址 ∩ msg.send_to ≠ ∅ → 命中(定向)

role.put_message(msg) → 塞进该角色私有信箱 role.py:448 / msg_buffer
┊ (只有命中的角色才收到)

▼ 下一轮 Environment.run() 时
role._observe(): 从信箱 pop_all,再按 cause_by/send_to 过滤出「我关心的」 role.py:399


该角色进入「观察→思考→行动」循环(见 01 章),产出新 Message,再次 publish……

谁在推着轮子转

Environment 只管「投递」。真正让所有角色一轮一轮跑起来的,是两层驱动:

  • Environment.run() —— 一轮之内,把所有非空闲角色用 asyncio.gather 并发跑一遍(base_env.py:197)。
  • Team.run(n_round) —— 外层大循环:每轮先查「全空闲了吗 / 钱够吗」,再调 env.run(),直到轮数耗尽或提前停(team.py:122)。

3. 核心机制(逐个拆解)

3.1 Message:路由三元组 + check_* 归一化

它要解决的小问题: 一封信得能回答三个问题——「因什么而发?谁发的?发给谁?」。这三个字段就是 MetaGPT 的路由三元组

字段常量名含义谁来读它
cause_byMESSAGE_ROUTE_CAUSE_BY因何而发:触发这封信的那个 Action 的类全名角色的 watch 订阅:「我盯着某类 Action 的产物」
sent_fromMESSAGE_ROUTE_FROM谁发的:发信角色的类全名调试 / 溯源
send_toMESSAGE_ROUTE_TO发给谁:一个收件人地址集合 set[str]Environment 分发时匹配

这三个常量定义在 metagpt/const.py:77-79。注意 send_to 默认是 {MESSAGE_ROUTE_TO_ALL},即 "<all>"——不指定收件人就是广播(schema.py:241)。

巧妙处:字段一律「归一化成字符串」。 你在业务里可以把 cause_by 写成一个 Action 类、把 send_to 写成一个角色对象或名字,但它们进入 Message 时会被一组 check_* 校验器统一转成字符串 / 字符串集合,这样分发时只做纯字符串比较,简单可靠。

# 真实源码,schema.py:266-279 —— 三个字段的归一化校验器(节选)
@field_validator("cause_by", mode="before")
def check_cause_by(cls, cause_by: Any) -> str:
# 不填就默认「用户需求」;任何 Action 类都被 any_to_str 转成类全名字符串
return any_to_str(cause_by if cause_by else import_class("UserRequirement", ...))

@field_validator("send_to", mode="before")
def check_send_to(cls, send_to: Any) -> set:
# 单个/列表/集合都被转成「字符串集合」,不填则回落到 {<all>}
return any_to_str_set(send_to if send_to else {MESSAGE_ROUTE_TO_ALL})

any_to_str(common.py:395)的规则很朴素:字符串原样返回,类 / 对象则取类全名(如 metagpt.actions.write_prd.WritePRD)。所以「订阅某个 Action」和「发给某个 Action 类名」能对得上。

除三元组外还有几个关键字段:id(check_iduuid4 自动补,schema.py:244)、role(system/user/assistant)、instruct_content(结构化产物,承载 02 章ActionNode 输出)、metadata

__setattr__ 的小陷阱(schema.py:307): 即使你事后赋值 msg.send_to = SomeRole,它也会被 __setattr__ 拦下来转成字符串集合——路由字段永远是干净的字符串,不会因为「后来手动改了一下」而破功。

便捷子类(schema.py:419-454): UserMessage / SystemMessage / AIMessage 只是预设了 roleMessage,方便对接 OpenAI 消息格式。其中 AIMessage 多了 with_agent(name) / agent 属性——把「这条 AI 回复出自哪个角色」记进 metadata,供 05 章 的 TeamLeader 调度识别来源。

3.2 两个路由「魔法地址」:<all><self>

const.py:81-83 定义了三个特殊地址:

常量作用
MESSAGE_ROUTE_TO_ALL"<all>"广播:任何角色都算收件人
MESSAGE_ROUTE_TO_NONE"<none>"谁都不发
MESSAGE_ROUTE_TO_SELF"<self>"发给我自己

<all>is_send_to 里被直接短路成「命中」(见 §3.3)。而 <self>在角色侧被翻译掉的——Role.publish_message(role.py:433-435)在发信前,把 send_to 里的 "<self>" 替换成发信角色自己的类全名:

# 真实源码,role.py:433-437 —— <self> 在离开角色时被就地翻译
if MESSAGE_ROUTE_TO_SELF in msg.send_to:
msg.send_to.add(any_to_str(self)) # 加上「我自己」的真实地址
msg.send_to.remove(MESSAGE_ROUTE_TO_SELF) # 摘掉占位符
if not msg.sent_from or msg.sent_from == MESSAGE_ROUTE_TO_SELF:
msg.sent_from = any_to_str(self) # 顺手补上「谁发的」

同一段还有个短路优化:如果这封信的收件人全是自己,角色直接 put_message 塞进自己信箱、根本不惊动环境(role.py:438-440)。只有要发给「别人」时,才真正 self.rc.env.publish_message(msg) 上总线。

3.3 Environment.publish_message:地址簿 + is_send_to 分发

它要解决的小问题: 一封信到了总线,怎么在一堆角色里挑出该收的那几个?

思路: 环境维护一本地址簿 member_addrs: Dict[BaseRole, Set](base_env.py:133)——每个角色对应它的一组地址。分发时逐个角色问 is_send_to:这封信的 send_to 和你的地址有交集吗?有就投。

# 真实源码,base_env.py:184-195 —— 分发的核心就这几行
def publish_message(self, message, peekable=True) -> bool:
found = False
for role, addrs in self.member_addrs.items(): # 遍历地址簿
if is_send_to(message, addrs): # 交集判定
role.put_message(message) # 命中 → 投进它的私有信箱
found = True
if not found:
logger.warning(f"Message no recipients: {message.dump()}")
self.history.add(message) # 全量留档(for debug)
return True

匹配逻辑 is_send_to(common.py:423-431)只有两条规则,先广播后定向:

def is_send_to(message, addresses: set):
if MESSAGE_ROUTE_TO_ALL in message.send_to: # ① 广播:含 <all> 直接命中
return True
for i in addresses: # ② 定向:地址与 send_to 有交集
if i in message.send_to:
return True
return False

两个值得记住的行为:

  • 没有收件人不会报错,只 warning,并且消息照样进 history——history 是一份全量调试留档(base_env.py:134 注释 # For debug),不参与路由。
  • put_message 只是把信 push 进角色的私有接收缓冲 rc.msg_buffer(role.py:448-452),并不立即处理。真正消费发生在下一轮该角色 _observe 时(见 §3.6 的说明与 01 章)。

3.4 地址簿从哪来:add_roles / set_addresses

它要解决的小问题: member_addrs 这本通讯录是谁、什么时候填的?

答案藏在「雇人」的链路里。Team.hireEnvironment.add_roles(base_env.py:164):

# 真实源码,base_env.py:164-173 —— 加入一批角色
def add_roles(self, roles):
for role in roles:
self.roles[role.name] = role # 先登记进 roles 名册
for role in roles:
role.context = self.context # 共享全局 Context
role.set_env(self) # 关键:反向把「环境」交给角色

role.set_env(env) 会回调 env.set_addresses(self, self.addresses)(role.py:308-313),把这个角色 → 它的地址集合写进 member_addrs(base_env.py:240-242)。那么一个角色的默认地址是什么?看 Role.check_addresses(role.py:211-213):

默认地址 = {角色类全名, 角色名},例如产品经理是 {"metagpt.roles.product_manager.ProductManager", "Alice"}

这解释了整套路由为什么能「对上」: 你只要把 send_to 写成对方的类名名字,就能命中它的地址集合。地址簿 = 角色注册时用「类名 + 名字」自动登记的双键索引。

set_addresses 也可被角色主动调用(role.py:293-300)来改自己的订阅地址,改完立刻同步回环境,地址簿始终最新。

3.5 Environment.run:并发跑非空闲角色 + is_idle + archive

它要解决的小问题: 一轮里,怎么让该干活的角色一起动、干完的别空转?

# 真实源码,base_env.py:197-211 —— 一轮驱动
async def run(self, k=1):
for _ in range(k):
futures = []
for role in self.roles.values():
if role.is_idle: # 空闲的跳过,不浪费
continue
futures.append(role.run()) # 收集协程
if futures:
await asyncio.gather(*futures) # 并发跑这一批

三个要点:

  • 并发而非串行。 一轮里所有非空闲角色用 asyncio.gather 同时推进;彼此靠上一轮投递到信箱的消息解耦,不需要显式排队。
  • is_idle 决定谁参与、以及全局是否停。 环境的 is_idle(base_env.py:228-234)是「所有角色都空闲」才为真;而单个角色的 is_idle(role.py:557-559)= 没有待处理新消息(news)、没有 todo、私有信箱 msg_buffer 也空。信箱一旦被投进新信,该角色就不再空闲,下一轮自动被拉起来。
  • archive:收尾存档(base_env.py:244-247)。 项目跑完时,若 context 里带 project_path,就对生成的代码仓做一次 GitRepository.archive()——把「一行需求跑成一个软件仓库」(04 章)的产物提交归档。

3.6 一封信从「投进信箱」到「被消费」

投递(§3.3)只把消息塞进 rc.msg_buffer,没有过滤。真正的「这封信我要不要理」发生在角色的 _observe(role.py:399-419):

# 真实源码,role.py:406-412(节选)—— 从信箱取信,再按订阅过滤
news = self.rc.msg_buffer.pop_all() # 清空信箱,取出所有新信
self.rc.news = [
n for n in news
if (n.cause_by in self.rc.watch # 我 watch 的 Action 触发的?
or self.name in n.send_to) # 或点名发给我的?
and n not in old_messages # 且没处理过(去重)
]

于是 MetaGPT 的路由其实是两级筛:

  1. 环境级(粗筛,§3.3): send_to ∩ 地址 —— 决定「信进不进你信箱」。
  2. 角色级(细筛,这里): cause_by ∈ watchname ∈ send_to —— 决定「进了信箱的信,你这轮理不理」。

cause_by 的价值在此凸显:角色不必知道会产出它要的东西,只需 watch 某类 Action——基于「因何而发」订阅,而非基于「谁发的」。这是 SOP 流水线(04 章)能自然串起来的底层原因。

3.7 Team:预算、轮次、存档

Team(team.py:32)是最外层的老板,给消息总线套上经营约束

雇人 hire(team.py:83): 就是转调 env.add_roles,把角色接进环境地址簿(§3.4)。

给钱 invest(team.py:92-96): 把投资额写进 CostManager.max_budget。运行时每轮 _check_balance(team.py:98-100)检查累计花费是否超预算,超了就抛 NoMoneyException——LLM 调用是要花钱的,这是硬性刹车

# 真实源码,team.py:98-100 —— 预算刹车
def _check_balance(self):
if self.cost_manager.total_cost >= self.cost_manager.max_budget:
raise NoMoneyException(self.cost_manager.total_cost, f"Insufficient funds: ...")

发起项目 run_project(team.py:102-107): 把一行需求包成一封普通 Message(content=idea)(默认广播、cause_by 归一为 UserRequirement)发上总线——这就是整条 SOP 的第一封信

主循环 run(team.py:122-138):

# 真实源码,team.py:128-137(节选)
while n_round > 0:
if self.env.is_idle: # 全员空闲 → 活干完了,提前收工
break
n_round -= 1
self._check_balance() # 钱不够 → 抛异常止损
await self.env.run() # 驱动一轮(§3.5)
self.env.archive(auto_archive) # 收尾归档
return self.env.history # 返回全量消息留档

停机有两个条件:轮数 n_round 耗尽,或 env.is_idle(没有任何消息在流动了)。@serialize_decorator(team.py:122)会在异常时自动把现场序列化,便于恢复。

存档 / 恢复(team.py:59-81): serialize 把整个 Team(含 env.context)写成 team.json;deserialize 反向重建——整台多智能体机器的状态可持久化、可断点续跑

3.8 ExtEnv:让环境也能承载「游戏 / 外部世界」

前面讲的 Environment 专注「角色间通信」。但它的父类 ExtEnv(base_env.py:53)还开了另一扇门:用 gym 式的读/写 API 把一个真实的外部环境(游戏、模拟世界)接进来,让角色像玩家一样 observe / step

机制是一对装饰器 + 两本注册表:

装饰器语义注册进
mark_as_readable(base_env.py:41)标记「从环境观察某状态」的方法env_read_api_registry
mark_as_writeable(base_env.py:47)标记「对环境施加某动作」的方法env_write_api_registry

被标记的方法会被登记进全局注册表;之后角色通过 read_from_api / write_thru_api(base_env.py:73 / :91)按名字调用它们——支持同步与协程两种实现。ExtEnv 还继承 gym 的 action_space / observation_space 和抽象的 reset / observe / step(base_env.py:106-121)。

一句话:Environment 把它当消息总线用,werewolf / minecraft / android 这类子环境把它当「可读可写的世界」用(见 EnvType,base_env.py:29)。二者共用同一套基类,是 MetaGPT「环境即可插拔载体」思想的体现。


4. 设计背景:RFC 113 / 116 的路由分层

源码里反复出现的 RFC 编号不是摆设,它记录了「为什么这么设计」。publish_message 的注释(base_env.py:176-183)直接引用了两份 RFC:

  • RFC 116 定的是消息路由的数据结构(即三元组 cause_by / sent_from / send_to):消息只声明收件人
  • RFC 113 定的是传输框架:「怎么把消息送到收件人」由传输层负责,消息本身不关心收件人在哪

这套「声明式路由 + 传输解耦」是本章所有代码的总纲:角色只说「我要发给谁 / 我 watch 什么」,至于跨进程、跨机器怎么送达,是 Environment(乃至未来更复杂的传输实现)的事。Team 存档时提到的 RFC 135(team.py:7-8)则补充了「项目完成后归档」这一步。


5. 边界与本章不覆盖的

  • 不覆盖 MGXEnv 的 TeamLeader 中转。 Team 默认 use_mgx=True,建的其实是 MGXEnv(team.py:43-51)——它在本章的「广播/定向」之上加了一个 TeamLeader(Mike,const.py:161)动态调度层。那套「谁该接下一棒由 TeamLeader 决定」的机制属于 05 章。本章讲的是底座:朴素 Environment 的总线语义。
  • history 不是路由的一部分。 它是全量留档,只服务调试 / 回放,别把它当消息队列用。
  • 投递「入信箱」≠「被处理」。 一条消息能否影响某角色,还要过角色级 _observewatch/send_to 细筛(§3.6);粗筛命中但细筛没命中的信,会静静躺在记忆里不触发行动。
  • publish_message 恒返回 True,即便没有任何收件人——「成功」只表示「分发流程跑完了」,不代表「有人收到」。要确认送达得看日志里的 Message no recipients 警告。

6. 代码地图(导航索引)

主题文件路径符号名
消息总线 / 环境本体metagpt/environment/base_env.pyEnvironment
分发核心:地址簿匹配投递metagpt/environment/base_env.py:175Environment.publish_message
一轮并发驱动非空闲角色metagpt/environment/base_env.py:197Environment.run
全局空闲判定(停机条件)metagpt/environment/base_env.py:228Environment.is_idle
雇人 / 登记地址簿metagpt/environment/base_env.py:164Environment.add_roles
写地址簿metagpt/environment/base_env.py:240Environment.set_addresses
项目收尾归档metagpt/environment/base_env.py:244Environment.archive
外部世界基类 + gym APImetagpt/environment/base_env.py:53ExtEnv / read_from_api / write_thru_api
读/写 API 注册装饰器metagpt/environment/base_env.py:41mark_as_readable / mark_as_writeable
收件人判定metagpt/utils/common.py:423is_send_to
路由字段字符串归一metagpt/utils/common.py:395any_to_str / any_to_str_set
消息 + 路由三元组metagpt/schema.py:232Message (cause_by/sent_from/send_to)
路由字段校验器metagpt/schema.py:266check_cause_by / check_send_to / check_sent_from
OpenAI 消息子类metagpt/schema.py:419UserMessage / SystemMessage / AIMessage
路由常量metagpt/const.py:77MESSAGE_ROUTE_TO_ALL / _TO_SELF / _CAUSE_BY
角色侧发布 + <self> 翻译metagpt/roles/role.py:429Role.publish_message
投进私有信箱metagpt/roles/role.py:448Role.put_message
信箱取信 + 订阅细筛metagpt/roles/role.py:399Role._observe
默认地址 = 类名+名字metagpt/roles/role.py:211Role.check_addresses
顶层老板:雇人/预算/轮次metagpt/team.py:32Team
主循环 + 停机条件metagpt/team.py:122Team.run
预算刹车metagpt/team.py:98Team._check_balance / NoMoneyException
存档 / 恢复metagpt/team.py:59Team.serialize / Team.deserialize

相邻章节: index.md(全景与阅读地图) · 01-role-loop.md(角色循环) · 02-action-actionnode.md(工作单元) · 04-classic-sop-pipeline.md(SOP 流水线) · 05-mgx-rolezero.md(MGX / TeamLeader 调度)