数据截至 (上游 commit 11cdf466d042)
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_by | MESSAGE_ROUTE_CAUSE_BY | 因何而发:触发这封信的那个 Action 的类全名 | 角色的 watch 订阅:「我盯着某类 Action 的产物」 |
sent_from | MESSAGE_ROUTE_FROM | 谁发的:发信角色的类全名 | 调试 / 溯源 |
send_to | MESSAGE_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_id 用 uuid4 自动补,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 只是预设了 role 的 Message,方便对接 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.hire → Environment.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 的路由其实是两级筛:
- 环境级(粗筛,§3.3):
send_to ∩ 地址—— 决定「信进不进你信箱」。 - 角色级(细筛,这里):
cause_by ∈ watch或name ∈ 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不是路由的一部分。 它是全量留档,只服务调试 / 回放,别把它当消息队列用。- 投递「入信箱」≠「被处理」。 一条消息能否影响某角色,还要过角色级
_observe的watch/send_to细筛(§3.6);粗筛命中但细筛没命中的信,会静静躺在记忆里不触发行动。 publish_message恒返回True,即便没有任何收件人——「成功」只表示「分发流程跑完了」,不代表「有人收到」。要确认送达得看日志里的Message no recipients警告。
6. 代码地图(导航索引)
| 主题 | 文件路径 | 符号名 |
|---|---|---|
| 消息总线 / 环境本体 | metagpt/environment/base_env.py | Environment |
| 分发核心:地址簿匹配投递 | metagpt/environment/base_env.py:175 | Environment.publish_message |
| 一轮并发驱动非空闲角色 | metagpt/environment/base_env.py:197 | Environment.run |
| 全局空闲判定(停机条件) | metagpt/environment/base_env.py:228 | Environment.is_idle |
| 雇人 / 登记地址簿 | metagpt/environment/base_env.py:164 | Environment.add_roles |
| 写地址簿 | metagpt/environment/base_env.py:240 | Environment.set_addresses |
| 项目收尾归档 | metagpt/environment/base_env.py:244 | Environment.archive |
| 外部世界基类 + gym API | metagpt/environment/base_env.py:53 | ExtEnv / read_from_api / write_thru_api |
| 读/写 API 注册装饰器 | metagpt/environment/base_env.py:41 | mark_as_readable / mark_as_writeable |
| 收件人判定 | metagpt/utils/common.py:423 | is_send_to |
| 路由字段字符串归一 | metagpt/utils/common.py:395 | any_to_str / any_to_str_set |
| 消息 + 路由三元组 | metagpt/schema.py:232 | Message (cause_by/sent_from/send_to) |
| 路由字段校验器 | metagpt/schema.py:266 | check_cause_by / check_send_to / check_sent_from |
| OpenAI 消息子类 | metagpt/schema.py:419 | UserMessage / SystemMessage / AIMessage |
| 路由常量 | metagpt/const.py:77 | MESSAGE_ROUTE_TO_ALL / _TO_SELF / _CAUSE_BY |
角色侧发布 + <self> 翻译 | metagpt/roles/role.py:429 | Role.publish_message |
| 投进私有信箱 | metagpt/roles/role.py:448 | Role.put_message |
| 信箱取信 + 订阅细筛 | metagpt/roles/role.py:399 | Role._observe |
| 默认地址 = 类名+名字 | metagpt/roles/role.py:211 | Role.check_addresses |
| 顶层老板:雇人/预算/轮次 | metagpt/team.py:32 | Team |
| 主循环 + 停机条件 | metagpt/team.py:122 | Team.run |
| 预算刹车 | metagpt/team.py:98 | Team._check_balance / NoMoneyException |
| 存档 / 恢复 | metagpt/team.py:59 | Team.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 调度)