Swarm runtime cn
AgentRuntime 是 Jianmu 多 Agent 协作的运行时引擎——它不像若干独立的单 Agent 循环那样各自为政,而是将消息路由、邮箱投递、事件总线和生命周期控制统一在同一个调度器之下。每当一个 Agent 通过 spawn() 诞生、通过 send_message() 发出消息、或在 auto 模式下持续 tick,背后都依赖本页所述的 Runtime 组件协同工作。
架构总览:从 AgentRuntime 看全局组件关系¶
AgentRuntime 是核心门面类(约 1160 行),其内部组合了六个关键子系统,各自承担明确职责。下图展示了它们的依赖方向与数据流:
flowchart TB
subgraph Runtime["AgentRuntime (core.py)"]
HOSTS["Host Registry<br/>_hosts: dict[str, AgentHost]"]
EB["Event Bus<br/>InMemoryEventBus"]
MS["Message Store<br/>InMemoryMessageStore"]
RIDX["Routing Indexes<br/>_group_members / _topic_subscribers / _attachments"]
TPS["Tool Providers<br/>[SwarmToolProvider, MCPToolProvider, ...user...]"]
PR["Prompt Runtime<br/>_prompt_runtime"]
end
subgraph Scheduler["scheduler.py"]
STEP["step_agent / step_all / run_until_idle"]
AUTO["run_runner_task / ensure_runner_task"]
RESTART["restart_policy + backoff"]
end
subgraph Mailbox["mailbox.py"]
SYNC["sync_messages"]
DRAIN["drain_mailbox"]
PRIME["prime_mailbox_for_run"]
GC["maybe_prune_message_store"]
end
subgraph Routing["routing.py"]
SEND["send_message"]
ENQ["enqueue_hot_envelope"]
WAKE["wake_targets"]
SUB["join_group / subscribe_topic"]
end
subgraph Host["AgentHost (host.py)"]
PROFILE["profile + role"]
RUNNER["ReactiveRunner"]
STATE["StateManager"]
INBOX["inbox: deque + pending_messages"]
end
Runtime --> Scheduler
Runtime --> Mailbox
Runtime --> Routing
HOSTS --> Host
Scheduler --> Host
Mailbox --> MS
Mailbox --> Host
Routing --> MS
Routing --> HOSTS
EB -->|"emit(AgentEvent)"| HOSTS
AgentRuntime.__init__ 接收约 20 个可配置参数,涵盖角色注册、消息存储、工具链、重启策略、快照检查点等全部维度。在构造阶段它会按顺序排列 ToolProvider 链:内置 SwarmToolProvider → 可选的 MCPToolProvider → 用户传入的额外 Provider,确保下游可覆盖上游同名工具。
AgentHost:每个 Agent 的运行时载体¶
AgentHost 是一个轻量级数据类(dataclass),它将单个 Agent 的所有运行时状态封装在一起,是调度器、邮箱和路由系统操作的原子单位。
| 字段 | 类型 | 用途 |
|---|---|---|
profile |
AgentProfile |
Agent 标识(id、role、task、parent_id) |
role |
AgentRole |
用于构建行为树的角色工厂 |
runner |
ReactiveRunner |
事件驱动的异步 tick 执行器 |
state |
StateManager |
类型化状态管理 |
mailbox_cursor_seq |
int |
已消费的消息序列号游标 |
groups |
set[str] |
当前订阅的广播组 |
topics |
set[str] |
当前订阅的主题 |
inbox |
deque[MessageEnvelope] |
热路径直达消息队列 |
pending_messages |
list[MessageEnvelope] |
待处理的消息列表 |
paused / auto_mode |
bool |
生命周期标志位 |
task / sync_task |
asyncio.Task |
后台 runner 任务和邮箱同步任务 |
mailbox_dirty / mailbox_force_signal / mailbox_version |
— | 邮箱变更检测机制 |
side_effects |
dict |
工具副作用去重(如 idempotent send/create) |
mailbox_version 是一个单调递增计数器,每次 mark_dirty() 调用都会使其递增。调度器在执行 runner 任务前后通过比较版本号来判断是否出现了新消息——这就是"新消息驱动继续运行"的检测机制。
消息路由:三种投递模式与热路径优化¶
消息从发送到被目标 Agent 消费,经历三个层次的处理管道:
sequenceDiagram
participant Sender as 发送方 Agent
participant Routing as routing.send_message
participant Store as InMemoryMessageStore
participant Inbox as 目标 Host.inbox
participant Mailbox as mailbox.sync_messages
participant State as StateMessageStore
Sender->>Routing: send_message(content, to_agent_id/group_id/topic)
Routing->>Store: append(MessageEnvelope) → seq
Routing->>Inbox: enqueue_hot_envelope(env)
Routing->>Sender: wake_targets()
Note over Sender: emit("message_sent")
Mailbox->>Inbox: 检查 inbox deque
Mailbox->>Store: read_since(cursor, agent_id, groups, topics)
Mailbox->>State: _append_many(incoming_messages)
Note over Mailbox: 可选 context_builder.transform 合并历史
Mailbox->>Store: ack(seq) 确认消费
三种路由目标¶
send_message 要求恰好指定一个目标(to_agent_id、group_id 或 topic),否则抛出 ValueError。三种模式的对比如下:
| 模式 | 目标参数 | 投递语义 | 适用场景 |
|---|---|---|---|
| 直连 (Direct) | to_agent_id |
精确送达指定 Agent | 点对点对话、任务委派 |
| 组播 (Group) | group_id |
送达组内所有成员 | 广播通知、团队同步 |
| 主题 (Topic) | topic |
送达所有订阅者 | 发布/订阅模式、事件驱动 |
热路径优化:enqueue_hot_envelope¶
消息存储到 InMemoryMessageStore 后,并不等到 Agent 下一次 mailbox sync 时才被发现——enqueue_hot_envelope 会直接将 MessageEnvelope 推入所有匹配目标 Agent 的 host.inbox deque。这意味着在同一个事件循环迭代中,目标 Agent 就能感知到新消息,大幅降低延迟。
运行时本地生产者必须使用 AgentRuntime.enqueue_internal_message(),而不是直接修改 host.inbox。这一统一交接会推进 mailbox version、记录可选的 mailbox_source,并可选地唤醒 Host,因此调度器不会遗漏本地注入的输入。
唤醒机制:wake_targets¶
消息发送后,wake_targets 负责通知接收方"有活干了"。直连模式下直接 runtime.wake(agent_id);组播/主题模式下遍历索引中的成员并逐一唤醒。它还支持 suppress_group_wake 元数据标记和显式 wake_agent_ids 列表,允许发送方精确控制唤醒范围。
组与主题管理¶
join_group / leave_group 和 subscribe_topic / unsubscribe_topic 操作 Agent 在 Runtime 级索引中的成员关系,同时发出对应的 group_joined、topic_subscribed 等事件。当 Agent 被 kill 时,remove_host_routing 会清理其在所有索引中的残留。
邮箱系统:消息同步与游标消费¶
邮箱系统在 mailbox.py 中实现了三个核心操作,它们共同构成了 Agent 消息消费的完整生命周期:
1. sync_messages:将外部消息载入 Agent 状态¶
这是消息消费的主入口。它按优先级从两个来源获取消息:
- 优先检查 host.inbox(热路径推送的消息),避免额外的存储查询
- 回退到 message_store.read_since(cursor) 从持久化存储读取
读取到的消息通过 StateMessageStore._append_many 写入 Agent 的状态管理器中的 messages 字段。如果 Runtime 配置了 context_builder,还会在合并历史消息时调用 context_builder.transform() 进行处理(例如 Token 预算控制、消息去重)。
每个成功同步的消息都会发出 message_received 事件,并推进 host.mailbox_cursor_seq 游标。
2. prime_mailbox_for_run:Runner 启动前同步¶
在每次 runner 执行前调用。它处理当前可能正在进行的 sync_task(等待其完成或清理),然后通过 sync_lock 获取锁进行同步。同步完成后,如果仍然 mailbox_dirty,会重新调度 schedule_mailbox_sync。
drain 到 prime 的交接¶
drain_mailbox 和 prime_mailbox_for_run 的执行顺序并不固定。若 drain 在新建 runner 到达 prime_mailbox_for_run 前已同步一批消息,会将该 incoming context 标记为已准备;即使 dirty 标志已清除,prime_mailbox_for_run 仍会为该次 execution 领取这批消息。未被 drain 明确标记为已准备的 context 则属于上一轮的陈旧上下文,会在新 execution 前清理。这样既避免吞掉新消息,也避免复用前一轮 execution 的 incoming batch。
3. drain_mailbox:有界迭代邮箱排空¶
在 auto 模式下,后台 sync_task 循环调用此函数。它在 mailbox_drain_max_iterations(默认 100)次迭代中反复检查 mailbox_dirty 标志并同步消息。每次迭代后都检查是否需要继续;达到上限时发出 mailbox_drain_capped 警告并调度后续继续。
死信处理¶
on_sync_failure 实现了消息投递的失败重试与死信队列机制。对于每次同步失败的 pending_messages,调用 message_store.fail() 累加失败计数;超过 delivery_max_failures(默认 3)阈值的消息通过 mark_dead() 移入死信队列,并发出 message_dead_letter 事件。
消息存储 GC¶
maybe_prune_message_store 在所有 Host 游标均已推进时触发。它计算所有 Host 游标的最小值作为安全删除边界,并确保两次 GC 之间的序列号跨度超过 message_store_gc_min_advance(默认 1000),以防止频繁回收。
InMemoryEventBus:线程安全的事件总线¶
InMemoryEventBus 是 Swarm 运行时内部事件的发布/订阅中枢。它基于环形缓冲区 + 订阅者列表设计,所有操作受 threading.Lock 保护。
classDiagram
class InMemoryEventBus {
-_subscribers: List[Callable]
-_buffer: List[Tuple[float, AgentEvent]]
-_max_buffer: int
-_lock: Lock
+emit(event) None
+subscribe(callback) Callable
+get_since(since) List[AgentEvent]
+get_all() List[AgentEvent]
}
class AgentEvent {
+event_type: str
+agent_id: str | None
+payload: dict
+created_at: float | None
+event_id: str | None
}
InMemoryEventBus --> AgentEvent
emit() 方法自动补全 event_id(UUID4)和 created_at(Unix 时间戳),然后将事件追加到缓冲区并通知所有订阅者。单个订阅者的异常不会影响其他订阅者。订阅通过返回 unsubscribe 回调来管理:调用该回调即从订阅列表中移除。
AgentRuntime.emit() 在调用 event_bus.emit() 之前会做一次规范化——将 session_id 和 agent_id 注入 payload——然后桥接选定的生命周期事件到 Runtime 级 RuntimeEventBus(agent.created、agent.started 等)。
调度器:手动模式与自动模式的双轨制¶
Swarm Runtime 提供两种 Agent 运行模式,分别对应不同场景:
| 维度 | 手动模式 (Manual) | 自动模式 (Auto) |
|---|---|---|
| 触发方式 | 显式调用 step_agent() 等 API |
spawn() 时 auto_start=True 或 start_agent() |
| 执行循环 | 调用方控制每步 tick | 后台 asyncio.Task 持续 runner.run() |
| 并发控制 | step_all_concurrent(max_concurrency=N) |
各自独立 task 并发 |
| 适用场景 | 调试、交互式步进、外部编排 | 全自动多 Agent 协作 |
手动模式 API¶
step_agent(agent_id):推进单个 Agent 一个 tick,返回Status。内部通过sync_lock完成消息同步后调用runner.step()。step_agent_with_actions(agent_id):同上,但额外返回actions字典。step_all():按主机注册顺序依次推进所有 Agent。step_all_concurrent(max_concurrency=N):通过asyncio.Semaphore并发推进所有 Agent。run_until_idle(max_steps=200):持续推进直到所有 Agent 处于空闲状态(无新消息、无 RUNNING 节点、无活跃 async_task)。
自动模式核心循环:run_runner_task¶
这是 auto 模式下每个 Host 的后台主循环,体现了 Runtime 的完整重启与唤醒语义:
flowchart TD
START["runner task 启动"] --> RESET["reset_runner_for_restart()"]
RESET --> PRIME["prime_mailbox_for_run()"]
PRIME --> RUN["runner.run()"]
RUN --> CHECK_RETRY{"收到消息但未发送?<br/>same_run_no_send_retries"}
CHECK_RETRY -->|是| RESET
CHECK_RETRY -->|否| CHECK_CONT{"新消息到达?<br/>should_continue_after_run"}
CHECK_CONT -->|是| RESET
CHECK_CONT -->|否| EMIT_DONE["emit agent_run_completed"]
RUN -->|异常| ERR_HANDLER{"restart_policy?"}
ERR_HANDLER -->|never| FAILED["emit agent_failed"]
ERR_HANDLER -->|on_failure/always| CHECK_RESTART{"restart_count < max_restarts?"}
CHECK_RESTART -->|是| BACKOFF["backoff 延迟"]
BACKOFF --> RESET
CHECK_RESTART -->|否| FAILED
should_continue_after_run 通过比较 mailbox_version 是否变化、检查 mailbox_dirty、以及调用 message_store.has_unread() 来判断运行期间是否有新消息到达,实现了消息驱动的连续运行。
should_retry_same_run_without_send 处理一种边缘情况:Agent 收到消息并完成了一轮运行,但在该轮中没有发送任何消息。此时如果 same_run_no_send_retries 允许,运行体会自动重试——确保"收到消息后一定回复"的语义。
空闲检测¶
is_idle 要求所有 Host 同时满足:mailbox_dirty 为 False、无活跃 sync_task、根节点状态非 RUNNING、且 has_active_work 返回 False。has_active_work 不仅检查消息队列,还遍历行为树中的所有节点:检查 has_pending_work 回调、活跃的 async_task、RUNNING 状态的 AsyncBehaviour。
Agent 生命周期:从创建到销毁的完整状态机¶
stateDiagram-v2
[*] --> Created: spawn(role, task)
Created --> Running: auto_start 或 start_agent()
Created --> Stepping: step_agent() (manual)
Running --> Paused: pause()
Paused --> Running: resume()
Running --> Preempted: preempt_agent()
Preempted --> Running: 自动 restart
Running --> Error: runner 异常
Error --> Running: restart (backoff)
Error --> Failed: restart 超限或 policy=never
Running --> Stopped: kill() / kill_async()
Stepping --> Stopped: kill() / kill_async()
Paused --> Stopped: kill() / kill_async()
Failed --> Stopped: kill() / kill_async()
Stopped --> [*]
spawn:从 Role 到 Host 的装配¶
spawn() 依次完成以下步骤:
1. 解析 role(字符串查找注册表或直接使用对象)
2. 生成 UUID agent_id,构建 AgentProfile
3. 处理 constraints(反序列化 → 预算规格提取 → 元数据注入)
4. 调用 role.build_tree(profile) 构建行为树根节点
5. 通过 state_factory 创建 StateManager 并初始化
6. 组装 RunContext(注入 model_client、runtime、sandbox、prompt_runtime)
7. 通过 runner_factory 创建 ReactiveRunner
8. 创建 AgentHost,处理 attached 模式父子关系
9. 注册到 _hosts 字典,发出 agent_created 事件
10. 若 auto_start=True,立即调用 start_agent()
kill:级联终止¶
kill() 通过 _detach_host 从 Runtime 中移除 Host,同时返回 attached 子 Agent 列表。然后对子 Agent 递归调用 kill(),实现级联终止。最后将分离的 sync_task 和 task 加入后台清理队列。
kill_async() 与之功能相同,但使用 await 等待任务取消完成,适用于需要确保清理完毕的异步上下文。
pause / resume / wake¶
pause:设置host.paused = True,发出agent_paused事件。已暂停的 Host 不再被调度器处理。resume:设置host.paused = False,发出agent_resumed,然后调用wake()触发工作。wake:标记 mailbox dirty + 强制信号 +ensure_runner_task,确保 Host 的 runner 任务存在并开始处理。
preempt_agent:中断并重启¶
preempt_agent 用于打断正在运行的 Agent 并强制重新开始。它取消当前 sync_task 和 task,通过 reset_runner_for_restart 重置 runner 状态(停止所有节点、取消异步任务),然后启动替代 runner task。替代 task 负责唯一一次 prime_mailbox_for_run,避免 mailbox drain 已准备的消息被重复领取或清除。
工具系统:Swarm 能力以 Tool 形式暴露¶
SwarmToolProvider 将 Runtime 的核心管理操作包装为 LLM 可调用的 Tool,通过 get_all_tools() 统一暴露:
| Tool 名称 | 功能 | 关键参数 |
|---|---|---|
create_agent / create |
创建子 Agent | role, guidance, reuse_key |
send_message |
发送路由消息 | to, content, group_id, topic, dedupe_key |
send_direct_message |
直连消息 | to_agent_id, content, dedupe_key |
send_topic_message |
主题广播 | topic, content, dedupe_key |
list_agents |
列出活动 Agent | — |
self |
返回当前 Agent 档案 | — |
pause_agent |
暂停 Agent | agent_id |
resume_agent |
恢复 Agent | agent_id |
kill_agent |
终止 Agent | agent_id |
这些工具内置幂等性支持:create_agent 使用 reuse_key 和 host.side_effects 字典确保相同 reuse_key 的重复调用返回已创建的 agent_id;send_message 使用 dedupe_key 和 effect_tags=("message_send",) 防止重复发送。
工具链采用 "后注册者覆盖同名工具" 的合并策略(ToolSet.from_tools(merged).with_tools(named_tools)),这使得用户可以注册自定义 Provider 来覆盖内置工具。
节点集成:SpawnAgent、SendMessage、WaitMessage¶
三个行为树节点将 Runtime 能力嵌入到 Agent 的执行树中:
-
SpawnAgent:异步节点。通过_resolve_runtime查找 Runtime 引用(支持属性注入、状态查找和全局回退),调用runtime.spawn()并将新的 agent_id 写入output_key。支持reuse_if_present模式避免重复创建。 -
SendMessage:异步节点。从状态中解析content、to_agent_id/group_id/topic和sender_id,调用runtime.send_message()。当 Runtime 不可用时降级为通过runtime.emit()发送事件。 -
WaitMessage:同步节点。仅检查ctx.incoming_messages是否非空——这是mailbox.sync_messages注入的上下文属性。有空消息时返回 SUCCESS 让行为树推进,否则返回 RUNNING 保持阻塞。
SwarmNode 继承自 SkillNode,增加了 swarm 特有的对话完成策略提示("必须调用 send_message 才算完成一轮"),以及通过 swarm_success 状态标志位记录执行结果。
快照与恢复:运行时持久化¶
save_snapshot 将整个 Runtime 状态序列化:
{
"version": 1,
"session_id": "...",
"hosts": [
{
"profile": {...}, # AgentProfile 字典
"state": {...}, # StateManager checkpoint
"mailbox_cursor_seq": N,
"groups": [...],
"topics": [...],
"attached_to": "...",
"paused": bool,
"auto_mode": bool,
...
}
],
"message_store": {...} # InMemoryMessageStore dump_state()
}
restore_snapshot 执行反向流程:先 kill 掉当前所有 Host → 恢复消息存储 → 逐个重建 Host(包括重新 build_tree、恢复状态、重建 RunContext 和 ReactiveRunner)→ 重建路由索引 → 重启 auto 模式任务。
序列化通过 CheckpointerProtocol 的 save_checkpoint/get_checkpoint 接口完成,当前支持 FileCheckpointer(文件系统后端)。
可观测性:事件流与桥接¶
Runtime 提供两层事件订阅:
-
subscribe_events(callback):订阅InMemoryEventBus上所有AgentEvent。每次事件在 publish 前经过emit()规范化(注入session_id和agent_id)。 -
subscribe_observability(callback, runtime_events=True, trace_events=False):统一的观测入口。runtime_events=True转发规范化的 Runtime 事件;trace_events=True额外订阅 telemetry trace 事件并以统一格式转发。
启用 trace 转发时,订阅会跟随进程级默认 TelemetryHub 的替换;telemetry bootstrap 或其他 Hub 替换后,调用方无需重新订阅。
Runtime 的完整事件类型清单涵盖:agent_created、agent_started、agent_restarted、agent_paused、agent_resumed、agent_wake、agent_preempted、agent_stopped、agent_killed、agent_failed、agent_error、agent_step_completed、agent_run_completed、agent_mailbox_continue、agent_no_send_after_incoming、agent_restart_scheduled、message_sent、message_received、message_dead_letter、mailbox_drained、mailbox_drain_capped、mailbox_sync_error、message_store_compacted、group_joined、group_left、topic_subscribed、topic_unsubscribed、runtime_restored。
其中选定的生命周期事件会桥接到 RuntimeEventBus(如 agent.created),使得上层 Agent 门面和遥测系统能够统一感知 Swarm 生命周期。
高层组装:SwarmRole 与 AgentTeam¶
SwarmRole(role.py)是面向"液态拓扑"协作场景的预设角色实现。它在 build_tree() 中构造一个 SwarmNode,内部采用 ReAct 循环 + 技能系统,将 system_prompt 和 task 拼接为完整 prompt,并通过 _resolve_runtime_tools() 获取 Runtime 工具。
AgentTeam(team.py)是 AgentRuntime 的外观包装器,负责:
- 从项目配置 jianmu.yaml 中读取 checkpoint 默认值
- 统一管理 defaults(model_client、context_builder、tool_providers)
- 提供 async with 上下文管理器接口(__aenter__ → initialize_async,__aexit__ → close_async)
两者在架构中的关系是:AgentTeam → AgentRuntime → AgentHost → ReactiveRunner(行为树)→ SwarmNode → SwarmRole(配置)。
阅读建议¶
本文档覆盖了 Swarm 运行时的内部引擎机制。接下来建议阅读:
- 角色与团队:AgentRole、SwarmRole 与 AgentTeam 的协作编排 —— 深入了解角色工厂、@role 装饰器和团队级配置
- Swarm 节点:SpawnAgent、SendMessage、WaitMessage 的多 Agent 通信原语 —— 从节点视角理解 Swarm 原语在行为树中的用法
- ReactiveRunner:事件驱动的异步 tick 调度与挂起恢复机制 —— AgentHost 底层的执行器详解