Swarm nodes cn
Swarm 节点是 Jianmu 多 Agent 协作的行为树原语集合,它将 AgentRuntime 的底层消息路由、Agent 生命周期管理和邮箱同步机制封装为可组合的 py_trees 节点。开发者通过在行为树中编排 SpawnAgent、SendMessage、WaitMessage 三个核心节点,即可构建出从简单的 Coordinator-Worker 协作到液态拓扑(liquid-topology)自组织网络在内的各类多 Agent 通信模式。这三个节点分别对应多 Agent 系统中最基础的三个动作:创建新 Agent、发送消息和等待接收消息——它们共同构成了一套完整的 CSP(Communicating Sequential Processes)风格的 Agent 间通信协议。
架构定位:节点如何接入 Swarm 运行时¶
在深入各节点实现之前,有必要理解它们在整体架构中的位置。Swarm 节点并不直接处理消息的存储、路由或投递——这些职责由 AgentRuntime 承担。每个 AgentRuntime 实例内部维护一个 InMemoryMessageStore(消息持久化)、三套路由索引(_group_members 组播、_topic_subscribers 主题订阅、直接点对点)以及一个 AgentHost 注册表。节点通过 _SwarmRuntimeNode 基类提供的运行时解析机制,在 update_async 时动态获取与之关联的 AgentRuntime 引用,随后调用其 spawn()、send_message() 等公开方法完成实际操作。
下面的 Mermaid 图展示了三个核心节点与 AgentRuntime 内部组件之间的调用关系:
graph TB
subgraph "行为树层 (Behavior Tree)"
SA[SpawnAgent]
SM[SendMessage]
WM[WaitMessage]
end
subgraph "AgentRuntime"
SP["spawn()"]
SND["send_message()"]
WK["wake()"]
MSG[InMemoryMessageStore]
HOST[AgentHost 注册表]
MBX[mailbox 同步]
end
subgraph "路由层 (Routing)"
RT[routing.send_message]
ENQ[enqueue_hot_envelope]
GRP[join_group / leave_group]
TOP[subscribe_topic]
end
SA -->|"runtime.spawn(role, task)"| SP
SP --> HOST
SM -->|"runtime.send_message(...)"| SND
SND --> RT
RT --> MSG
RT --> ENQ
ENQ --> HOST
ENQ --> WK
WK --> MBX
MBX -->|"写入 ctx.incoming_messages"| WM
WM -->|"检查 ctx.incoming_messages"| MBX
关键点在于:SendMessage 将消息写入消息存储并通过 enqueue_hot_envelope 直接推送到目标 AgentHost 的 inbox 队列,随后 wake() 触发目标 Agent 的邮箱同步;WaitMessage 则通过检查 ctx.incoming_messages 来感知新消息的到达——这个属性正是由 mailbox 同步过程写入的。
_SwarmRuntimeNode:运行时解析基类¶
三个 Swarm 节点都直接或间接继承自 _SwarmRuntimeNode,这个内部基类继承 AsyncNode,提供了两个关键的辅助方法:
_resolve_runtime() 采用三级回退策略查找 AgentRuntime 实例:优先使用构造时传入的 runtime 属性;若为 None,从状态管理器中读取 "runtime" 键(支持全局回退);最后尝试从整个状态对象上通过 getattr(state, "runtime") 获取。解析后还会验证运行时是否具备 required_methods 中指定的方法——例如 SpawnAgent 要求 "spawn",SendMessage 要求 "send_message" 和 "emit"。
_read_state_value() 封装了带命名空间解析的状态读取逻辑,先通过 resolve_input_port 解析键名,然后在端口绑定的状态键上读取;若未找到且 allow_global_fallback=True,还会在全局状态中查找。
这种设计使得节点既可以在构造时显式绑定运行时(适用于静态流程),也可以依赖状态管理器在运行时动态获取(适用于 SwarmRole 等模板化角色场景)。
SpawnAgent:Agent 生成原语¶
SpawnAgent 是动态创建 Agent 的行为树节点。它调用 AgentRuntime.spawn() 完成实际的 Host 创建、行为树构建和状态初始化,并将新生成的 agent_id 写入状态以供后续节点使用。
参数体系¶
| 参数 | 类型 | 默认值 | 说明 |
|---|---|---|---|
role |
str \| None |
None |
静态角色名,设置后忽略 role_key 查找 |
role_key |
str |
"role" |
从状态中读取角色名的键 |
task |
str \| None |
None |
静态任务描述,设置后忽略 task_key 查找 |
task_key |
str |
"task" |
从状态中读取任务的键 |
parent_id |
str \| None |
None |
静态父 Agent ID,设置后忽略 parent_id_key |
parent_id_key |
str |
"agent_id" |
从状态中读取父 Agent ID 的键 |
mode |
str |
"detached" |
生成模式:"detached"(独立)或 "attached"(依附父 Agent) |
constraints |
object \| None |
None |
静态执行约束,覆盖 constraints_key 查找 |
constraints_key |
str |
"constraints" |
从状态中读取约束的键 |
output_key |
str |
"spawned_agent_id" |
新 Agent ID 写入状态的键 |
reuse_if_present |
bool |
False |
若为 True,当 output_key 已有值时复用而不重新生成 |
执行逻辑¶
update_async() 首先解析运行时引用——若运行时不存在或不具备 spawn 方法,直接返回 FAILURE。接着检查 reuse_if_present 模式:如果启用了复用且 output_key 对应的状态槽中已存在非空字符串,则维持该值并返回 SUCCESS,实现幂等的 Agent 创建语义。
在正常生成路径中,节点依次解析 role(必填,缺失则 FAILURE)、task、parent_id 和 constraints,随后调用 runtime.spawn(role=..., task=..., parent_id=..., mode=..., constraints=...)。AgentRuntime.spawn() 内部执行以下步骤:
- 解析角色对象并生成 UUID 作为
agent_id - 构建
AgentProfile并调用角色的build_tree(profile)构建行为树根节点 - 通过
state_factory创建StateManager并注入初始状态(agent_id、role、task) - 组装
RunContext(包含 model_client、runtime、constraints、sandbox 等) - 创建
AgentHost并注册到_hosts字典 - 若
auto_start=True,立即启动该 Agent 的后台运行任务
生成的 agent_id 通过 write_port("spawned_agent_id", agent_id, state_key=self.output_key) 写入状态,后续节点通过读取该键获取目标 Agent 标识。
attached 模式与父 Agent 生命周期¶
当 mode="attached" 且 parent_id 非空时,子 Agent 会被记录到父 Agent 的 _attachments 集合中。这意味着当父 Agent 被 kill 时,所有 attached 子 Agent 也会被级联终止。"detached" 模式下的 Agent 则拥有独立的生命周期。
SendMessage:消息发送原语¶
SendMessage 是向其他 Agent、组或主题发送消息的行为树节点。它封装了 AgentRuntime.send_message() 的完整路由语义,支持三种目标寻址方式。
参数体系¶
| 参数 | 类型 | 默认值 | 说明 |
|---|---|---|---|
to_agent_id |
str \| None |
None |
静态直接消息目标 Agent ID |
to_agent_key |
str |
"to_agent_id" |
从状态读取目标 Agent ID 的键 |
group_id |
str \| None |
None |
静态组播目标组 ID |
group_key |
str |
"group_id" |
从状态读取组 ID 的键 |
topic |
str \| None |
None |
静态主题名称 |
topic_key |
str |
"topic" |
从状态读取主题的键 |
content |
str \| None |
None |
静态消息正文,设置后忽略 content_key |
content_key |
str |
"message" |
从状态读取消息正文的键 |
sender_id |
str \| None |
None |
静态发送者 ID,设置后忽略 sender_key |
sender_key |
str |
"agent_id" |
从状态读取发送者 ID 的键 |
路由语义¶
SendMessage 要求 to_agent_id、group_id、topic 三者恰好有一个非空——这是由底层 routing.send_message() 强制的约束:若零个或多个目标被同时指定,将抛出 ValueError。这种严格的路由排他性避免了消息被意外广播或路由歧义。
双路径发送策略¶
update_async() 实现了两条发送路径,体现了渐进增强的设计理念:
主路径(runtime.send_message):当解析到的运行时具备 send_message 方法时(AgentRuntime 实例具备),直接调用 await runtime.send_message(...),消息经过完整的存储持久化 → 信封入队 → 目标唤醒流程。
回退路径(事件发射):当运行时仅具备 emit 方法但不具备 send_message 时(例如某些精简运行时实现),节点自行构造 Message 对象并封装为 AgentEvent(event_type="message_sent", ...) 通过 runtime.emit() 发布。这条路径确保节点在最小运行时契约下仍可工作,但不会触发自动的消息存储和目标唤醒。
这种双路径策略使 SendMessage 保持了与 AgentRuntimeProtocol 中不同实现层级的兼容性。
消息投递的完整链路¶
当使用 AgentRuntime 的主路径时,一次消息发送的完整链路如下:
sequenceDiagram
participant SM as SendMessage
participant RT as AgentRuntime
participant Routing as routing.send_message
participant Store as InMemoryMessageStore
participant Enq as enqueue_hot_envelope
participant Host as 目标 AgentHost
participant MBX as mailbox.sync_messages
participant WM as WaitMessage (目标)
SM->>RT: send_message(sender_id, content, to_agent_id)
RT->>Routing: send_message(...)
Routing->>Store: append(envelope) → seq
Routing->>Enq: enqueue_hot_envelope(env)
Enq->>Host: host.inbox.append(envelope)
Enq->>Host: mark_dirty + wake
Routing->>RT: emit("message_sent")
RT->>MBX: schedule mailbox sync
MBX->>Host: read inbox → set ctx.incoming_messages
WM->>WM: 检查 ctx.incoming_messages → SUCCESS
WaitMessage:消息等待原语¶
WaitMessage 是所有 Swarm 节点中唯一继承同步 Node(而非 AsyncNode)的节点。它不执行异步 I/O,而是通过轮询 ctx.incoming_messages 属性来判断是否有新消息到达。
设计原理¶
WaitMessage 的同步特性源于 mailbox 系统的前置工作模式:在每次 Agent 的 tick 周期开始前,调度器(scheduler.step_agent)会先调用 runtime._sync_messages(host) 完成邮箱同步,将待投递消息写入 ctx.incoming_messages 和 ctx.incoming_envelopes。因此,当 WaitMessage.update() 被调用时,ctx.incoming_messages 已经是同步后的最新状态,无需额外的异步操作。
行为语义¶
def update(self) -> Status:
incoming = getattr(self.ctx, "incoming_messages", None)
if isinstance(incoming, list) and incoming:
return Status.SUCCESS
return Status.RUNNING
节点的逻辑极其简洁:有消息则 SUCCESS,无消息则 RUNNING。这种设计使其天然适配行为树的记忆序列(memory=True 的 Sequence):当 WaitMessage 返回 RUNNING 时,父序列会挂起等待,下一次 tick 时重新从该节点开始执行,直到消息到达后继续执行序列中的后续子节点。
ctx.incoming_messages 的数据来源¶
mailbox.sync_messages() 在同步消息时不仅设置 ctx.incoming_messages(消息对象列表),还设置 ctx.incoming_envelopes(信封对象列表,包含 seq、路由元数据等),同时将消息合并到状态管理器的 "messages" 键中(通过 StateMessageStore._append_many)。这使得下游节点不仅可以通过 WaitMessage 感知消息到达,还可以直接访问 ctx.incoming_messages 获取消息内容、发送者等详细信息。
WaitForever:永不终止的辅助节点¶
WaitForever 是 WaitMessage 的镜像节点——它始终返回 Status.RUNNING,永不自行终止。其主要用途是让一个 Agent 在完成初始化工作后保持存活状态,等待外部中断(如 preempt_agent 或 kill)。在需要 Agent 长期驻留等待动态消息的场景中,WaitForever 常被放在行为树末尾,确保 Agent 不会因树执行完毕而被标记为完成。
SwarmNode:基于技能的 Swarm Agent 节点¶
SwarmNode 继承自 SkillNode,是为 SwarmRole 定制的技能执行节点。它在 SkillNode 的基础上做了三项关键调整:
输出 Schema 抑制:_resolve_output_schema() 覆盖为返回 None(除非显式设置),因为 Swarm Agent 的最终输出通常是通过 send_message / send_direct_message 发送的纯文本,而非技能的 Pydantic 结构化输出。
输出键策略:_resolve_output_key() 优先使用显式设置的 output_key,否则使用节点名称——这避免了单技能场景下自动继承技能专用输出键的行为。
对话完成策略注入:_skill_react_system_prompt() 在基础 Prompt 后追加了"对话完成策略"指令,明确告知 LLM:一次完整的 turn 必须在调用 send_message 或 send_direct_message 发送可见回复后才算完成,避免 Agent 在内部推理或工具执行后提前终止。
此外,update() 在执行父类逻辑后,会将 swarm_success 标记写入状态管理器,供外部观察者(如 AgentTeam 或测试代码)判断 Agent 是否成功完成。
完整通信模式:Coordinator-Worker 协作¶
以下示例来自 basic_collaboration_demo.py,展示了三个核心节点协同工作的最小完整模式:
sequenceDiagram
participant C as Coordinator Tree
participant R as AgentRuntime
participant W as Worker Tree
C->>R: SpawnAgent("SpawnWorker", role="worker")
R->>W: 创建 AgentHost + 启动
R-->>C: agent_id → state["worker_id"]
C->>C: PrepareTaskNode: 构造 task 消息
C->>R: SendMessage(to_agent_key="worker_id", content_key="task")
R->>W: 投递消息到 Worker inbox
R->>W: mailbox sync → ctx.incoming_messages
W->>W: WaitMessage → SUCCESS (有消息)
W->>W: BuildWorkerReplyNode: 构造回复
W->>R: SendMessage(to_agent_key="reply_to", content_key="reply")
R->>C: 投递回复到 Coordinator inbox
C->>C: WaitMessage → SUCCESS (有回复)
C->>C: CaptureReplyNode: 提取 final_reply
Coordinator 行为树是一个带记忆的 Sequence:
Sequence(name="CoordinatorFlow", memory=True, children=[
SpawnAgent("SpawnWorker", role="worker", output_key="worker_id",
runtime=self.runtime, reuse_if_present=True),
PrepareTaskNode(),
SendMessage("SendTaskToWorker", to_agent_key="worker_id",
content_key="task", runtime=self.runtime),
WaitMessage("WaitWorkerReply"),
CaptureReplyNode(),
])
Worker 行为树同样使用记忆序列:
Sequence(name="WorkerFlow", memory=True, children=[
WaitMessage("WaitTask"),
BuildWorkerReplyNode(),
SendMessage("SendReply", to_agent_key="reply_to",
content_key="reply", runtime=self.runtime),
])
reuse_if_present=True 的设置在 Coordinator 的行为树中至关重要:若行为树因某种原因被重新执行(例如记忆序列的回退),已有的 worker_id 不会被覆盖,从而避免创建冗余 Agent。
与 Swarm 运行时工具的互补关系¶
三个核心节点面向静态行为树编排——开发者在设计时明确 Agent 之间的通信拓扑。与之互补的是 SwarmToolProvider 提供的运行时工具集,它们允许 LLM 驱动的 Agent(通过 SwarmRole / SwarmNode)在运行时动态做出通信决策:
| 工具名 | 对应节点 | 差异 |
|---|---|---|
create / create_agent |
SpawnAgent |
LLM 动态决定何时创建、创建什么角色;支持 reuse_key 幂等 |
send_message / send_direct_message / send_topic_message |
SendMessage |
LLM 选择目标、内容;支持 dedupe_key 幂等 |
subscribe_topic |
无直接节点对应 | 运行时动态加入主题,WaitMessage 随后可接收该主题消息 |
list_agents / self / pause_agent / resume_agent / kill_agent |
无直接节点对应 | 提供 Agent 自省和生命周期管理能力 |
这种"静态节点 + 动态工具"的双层设计使 Jianmu 的多 Agent 系统既能支持确定性的协作流程(如 Coordinator-Worker),也能支持 LLM 自组织的液态拓扑(如 autonomous_swarm_demo.py 中的 Leader-Scout-Analyst 三级链式路由)。
阅读建议¶
至此,你已经掌握了构建多 Agent 通信的节点原语。要理解这些节点背后的完整运行时机制,建议继续阅读:
- Swarm 运行时:AgentRuntime 的邮箱路由、事件总线与生命周期管理 — 深入了解
AgentRuntime.spawn()和send_message()的底层实现,包括消息存储、路由索引和 Agent Host 的生命周期。 - 角色与团队:AgentRole、SwarmRole 与 AgentTeam 的协作编排 — 了解如何使用
SwarmRole和AgentTeam外观模式组装 LLM 驱动的动态 Agent,以及它们如何利用 Swarm 工具集实现自组织协作。 - 行为树执行内核:基于 py_trees 的异步扩展与 Jianmu 节点模型 — 理解
AsyncNode、Node基类以及记忆序列(memory=True)如何支撑WaitMessage的挂起-恢复语义。