1 张图 + 6 阶段流水线:教你从零拆解 AI Agent 端到端架构(nanobot 源码)
1 张图 + 6 阶段流水线:教你从零拆解 AI Agent 端到端架构(nanobot 源码)
本文回答什么问题:如何设计一个端到端的 AI Agent 架构,让消息从聊天通道进入后能在 5 层管道中稳定流转并最终回到用户?
目标读者:LLM Agent 开发者、系统架构师
预计阅读时间:12 分钟
源码版本:nanobot v2025.10(GitHub HKUDS/nanobot)
nanobot 是一款轻量级 AI Agent 框架(30+ 核心模块、50+ Provider、16+ 通道),本文基于其源码,从一次"用户在 Telegram 发送消息"的真实场景出发,拆解端到端 5 层管道架构与 AgentLoop 内部的 6 阶段流水线。读完后你能画出自己的 AI Agent 架构图。
1.1 一次完整对话的全景图
下面是一张从"用户在 Telegram 发送消息"到"Agent 通过 Telegram 给出响应"的端到端架构图:
关键路径解读:
- 入站路径(蓝色流程):Telegram 客户端 → TelegramChannel 适配器 → MessageBus.inbound → AgentLoop.run() 1s tick
- 编排路径:AgentLoop →
_dispatch(SessionLock 串行同 session / ConcurrencyGate 限制跨 session 并发) →_process_message六阶段流水线 - 执行路径:AgentRunner 的
for iteration循环 →_request_model调 LLM → 流式 chunks 通过CompositeHook翻译成OutboundMessage→ 工具调用走_execute_tools→ 注入消费走_try_drain_injections - Provider 路径:
LLMProvider抽象 →AnthropicProvider/OpenAICompatProvider/FallbackProvider(装饰器模式) - 出站路径:
TurnDelivery投递 → MessageBus.outbound →ChannelManager._dispatch_outbound按事件类型分发(合并 / 抑制 / retry) → channel.send() → 回到 Telegram 客户端
1.2 关键路径详解
1.2.1 入站路径:Channel → AgentLoop
InboundMessage 定义在 [nanobot/bus/events.py
@dataclass
class InboundMessage:
channel: str
sender_id: str
chat_id: str
content: str
media: list[dict] = field(default_factory=list)
metadata: dict[str, Any] = field(default_factory=dict)
session_key_override: str | None = None
timestamp: float = field(default_factory=time.time)
@property
def session_key(self) -> str:
if self.session_key_override:
return self.session_key_override
return f"{self.channel}:{self.chat_id}"
1.2.2 编排路径:AgentLoop 接管
AgentLoop.run()([nanobot/agent/loop.py](file:///e:/learn/nanobot/nanobot/agent/loop.py#L1079-L1170))每 1 秒醒一次:
async def run(self) -> None:
while not self._shutdown.is_set():
try:
msg = await asyncio.wait_for(self.bus.consume_inbound(), timeout=1.0)
except asyncio.TimeoutError:
await self._check_expired_sessions_if_due()
continue
await self._dispatch(msg)
_dispatch() 根据 session 锁串行化同一会话、不同 session 并发:
async def _dispatch(self, msg: InboundMessage) -> None:
session_key = self._resolve_session_key(msg)
async with self._session_locks.setdefault(session_key, asyncio.Lock()):
async with self._concurrency_gate: # 默认上限 3
await self._process_message(msg, session_key)
_process_message 走完整的六阶段流水线:
async def _process_message(self, msg, session_key):
ctx = await self._restore_turn(msg, session_key) # restore
session, summary = await self._compact_session(ctx) # compact
if await self._dispatch_command(msg, ctx): # command
return # 斜杠命令已处理
await self._build_turn(ctx, session, summary) # build
response = await self._run_turn(ctx) # run
await self._persist_turn(ctx, response) # save
await self._prepare_outbound(ctx, response) # respond
1.2.3 执行路径:AgentRunner 调用 LLM
AgentRunner.run()([nanobot/agent/runner.py])是纯 LLM 循环:
async def run(self, spec: AgentRunSpec) -> AgentRunResult:
messages = list(spec.initial_messages)
for iteration in range(spec.max_iterations):
# 1. 注入检查
should_continue, _ = await self._try_drain_injections(spec, messages)
if not should_continue:
break
# 2. 调用 LLM
response = await self._request_model(spec, messages)
# 3. 落地 assistant message
messages.append(self._build_assistant_message(response))
# 4. 工具调用
if response.tool_calls:
tool_results = await self._execute_tools(spec, response.tool_calls)
messages.extend(tool_results)
continue
# 5. 没有工具调用 = 终止
return AgentRunResult(final_content=response.content, ...)
return self._max_iterations_fallback(spec, messages)
1.2.4 出站路径:Hook → Delivery → Channel
每一次 LLM 增量、工具完成、思考结束,都通过 CompositeHook 翻译成 OutboundMessage,注入到 bus.outbound:
核心要点速查(建议收藏)
- nanobot 采用 5 层管道架构:Channel → Bus → AgentLoop → AgentRunner → Provider([nanobot/agent/loop.py])
- 6 阶段流水线:restore → compact → command → build → run → save → respond,每阶段独立计时/插桩
- 关键设计:1s tick 主循环([loop.py L1079])兼顾巡检,两层锁(SessionLock + ConcurrencyGate Semaphore max=3)
- 数据流:InboundMessage([bus/events.py L23-L58])入站 → OutboundMessage(联合类型含 5 种 event)出站,全程流式
- 崩溃恢复:Checkpoint 自动持久化(
_set_runtime_checkpoint),OOM 后下次启动自动续传,最多 12 轮 - 出站分 4 步:CompositeHook 流式翻译 → TurnDelivery 包装 → MessageBus.outbound → ChannelManager 按 event 类型分发
ChannelManager._dispatch_outbound([nanobot/channels/manager.py])按 event 类型分发:
async def _dispatch_outbound(self, msg: OutboundMessage) -> None:
event = msg.event
if isinstance(event, ProgressEvent):
if event.reasoning_delta: ...
if event.tool_hint: ...
if event.file_edit_events: ...
elif isinstance(event, StreamDeltaEvent):
await self._coalesce_stream_deltas(msg)
await channel.send_delta(msg)
elif isinstance(event, StreamEndEvent):
await channel.send_delta_complete(msg)
elif isinstance(event, StreamedResponseEvent):
await channel.send(msg)
# ... 重试包装
await self._send_with_retry(channel, msg)
1.3 上下文注入机制
AgentLoop 用 contextvars 把每个 turn 的运行时状态注入到任意深度的调用栈:
四个关键 ContextVar:
# nanobot/agent/tools/context.py
_CURRENT_REQUEST_CONTEXT: ContextVar[RequestContext | None] = ContextVar(
"request_context", default=None
)
# nanobot/security/workspace_access.py
_CURRENT_WORKSPACE_SCOPE: ContextVar[WorkspaceScope | None] = ContextVar(
"workspace_scope", default=None
)
# nanobot/agent/tools/file_state.py
_CURRENT_FILE_STATES: ContextVar[FileStates | None] = ContextVar(
"file_states", default=None
)
# nanobot/agent/memory.py
_DREAM_CURSOR: ContextVar[int | None] = ContextVar(
"dream_cursor", default=None
)
工具可以直接读取:
# 工具内任意位置
ctx = current_request_context()
print(ctx.channel, ctx.chat_id, ctx.session_key)
这避免了把所有上下文作为参数传遍整个调用栈(典型的 request: Request = ... 污染)。
1.4 检查点与崩溃恢复
AgentLoop 在 session metadata 中保存:
# nanobot/agent/loop.py
_RUNTIME_CHECKPOINT_KEY = "_runtime_checkpoint"
_PENDING_USER_TURN_KEY = "_pending_user_turn"
Checkpoint 场景:
# 工具执行后
self._set_runtime_checkpoint(
session_key,
assistant_message=last_assistant,
completed_tool_results=results,
pending_tool_calls=next_batch,
)
恢复时机:
# _restore_turn 时
checkpoint = session.metadata.get(_RUNTIME_CHECKPOINT_KEY)
if checkpoint:
# 把已完成的结果回填到 messages
# 然后重放 pending_tool_calls
Pending User Turn:
# 用户消息已持久化,但 assistant 还没生成 → 崩溃
# 下次启动时检测到 → 补全 assistant 占位 "[Previous assistant message omitted.]"
1.5 持久化原子写
凡是"可能写一半"的状态都用 tmp + fsync + os.replace 模式:
# nanobot/agent/memory.py: MemoryStore.append_history
def append_history(self, entry: str, session_key: str = ...) -> int:
with self._append_lock:
path = self._history_path()
tmp = path.with_suffix(path.suffix + ".tmp")
with open(tmp, "a", encoding="utf-8") as f:
f.write(line)
f.flush()
os.fsync(f.fileno())
os.replace(tmp, path) # 原子替换
return cursor
类似地,[nanobot/session/manager.py]) 用 * .tmp + os.replace;[nanobot/cron/service.py]用 _atomic_write + 父目录 fsync;[nanobot/pairing/store.py] 用 _write_text_atomic。
1.6 死循环防护
下面三层防护让 nanobot 不会因为 LLM 抽风而无限循环:
max_iterations(默认 200)由AgentLoop注入到AgentRunSpec.max_iterations,AgentRunner跑到上限就强制退出_MAX_INJECTION_CYCLES = 5在 [nanobot/agent/runner.py]限制单回合最多 5 次注入_MAX_GOAL_CONTINUATION_ROUNDS = 12在 [nanobot/session/turn_continuation.py] 限制持续目标最多 12 轮续传
1.7 并发模型
三层并发语义:
- 同一 session_key 串行(
_session_locks) - 跨 session 受
ConcurrencyGate(Semaphore,默认 3)限制 - 同 session 内中途注入走
_pending_queues[key](maxsize=20),避免抢占
1.8 配置热更新
watchfiles.awatch([nanobot/config/watcher.py])仅盯目标文件:
async def watch_config_file(path: Path, on_change: Callable[[], None]) -> None:
parent = path.parent.resolve()
target = path.resolve()
async for changes in awatch(parent):
for _change_type, changed_path in changes:
if Path(changed_path).resolve() == target:
on_change()
break
回调链路:
1.9 启动顺序
nanobot gateway (CLI)
↓
_load_runtime_config → resolve_config_env_vars
↓
ChannelManager.discover_plugins
↓
_provider_snapshot_loader(config) → make_provider + FallbackProvider
↓
Consolidator + MemoryStore + SessionManager
↓
AgentLoop.from_config(provider_snapshot_loader, ...)
↓
CronService + LocalTriggerStore
↓
register cron callbacks (dream / heartbeat)
↓
gateway_service.start() → main loop
1.10 终止与清理
SIGINT / SIGTERM 二次按下触发强制 shutdown:
# nanobot/cli/commands.py
def _install_signal_handlers(loop, shutdown_event):
# 第一次:设置 shutdown_event + 触发 cron.stop()
# 第二次:取消所有 task → agent.sessions.flush_all()
清理顺序:
1.11 关键设计权衡
权衡 1:消息总线的极简性
MessageBus([nanobot/bus/queue.py])只 32 行,只持两个 asyncio.Queue。所有语义(流式、进度、思考、retry)都通过 OutboundMessage.event 字段承担。
好处:极易测试;新通道只需 publish_inbound + consume_outbound_consume。
代价:bus 本身不能做"按事件路由到不同 channel"——必须由 ChannelManager._dispatch_outbound 自己做。
权衡 2:Provider 抽象的"宽进严出"
LLMProvider.chat() / chat_stream() 是抽象。AnthropicProvider 直接用 AsyncAnthropic SDK;OpenAICompatProvider 用 AsyncOpenAI SDK 处理 30+ provider。
好处:统一重试 / fallback / 流式语义。
代价:必须满足"相同的最小字段"——但现实是各家 schema 差异大,导致 OpenAICompatProvider 1750 行。
权衡 3:Tool 自发现 vs 显式注册
ToolLoader 用 pkgutil.iter_modules 扫描 nanobot/agent/tools/ 下所有模块。但工具作者必须继承 Tool 并实现 execute。
好处:新工具"装完即用"。
代价:测试时容易意外加载所有工具——所以 _SKIP_MODULES 必须明确列出。
权衡 4:持久化选 JSONL + git
MemoryStore 把 MEMORY.md / USER.md / SOUL.md 放在 git 仓库里管理;history.jsonl 走 append-only;session.jsonl 走 message snapshot。
好处:审计可重现(git log)、人类可读、跨平台。
代价:大 session 时读写放大;并发写需要文件锁。
1.12 关键架构决策盘点
| 决策 | 选择 | 原因 |
|---|---|---|
| 异步框架 | 全 asyncio | 与 LLM 调用的 I/O 密集型匹配 |
| HTTP 客户端 | 异步 httpx | 高并发 + 流式 |
| 类型系统 | Pydantic v2 | 配置 / 事件 / 内部模型 |
| 序列化 | JSON + JSONL | 人类可读 + 跨平台 |
| 锁机制 | asyncio.Lock + contextvars |
异步上下文 |
| 并发控制 | Semaphore + asyncio.Lock 两层 |
全局 + 局部 |
| Provider 适配 | 抽象 + fallback 装饰器 | 业务无需感知 |
| Tool 自发现 | pkgutil + entry_points |
零侵入 |
| Channel 自发现 | 包扫描 + 懒加载 | 缺失依赖不阻塞 |
| 记忆整合 | append + 摘要 + 离线 | 减少 token 压力 |
| 错误恢复 | 重试 + circuit breaker + 流恢复 | 兼容网络抖动 |
| 安全模型 | 三层 scope (路径 / 网络 / 资源) | 深度防御 |
| 配置 | JSON + 环境变量 + 热更新 | 灵活 + 易部署 |
| 凭据管理 | OAuth 单独存 + prox env | 不入 config.json |
| WebUI 协议 | WebSocket 多路复用 | 实时 + 多端 |
| API 协议 | OpenAI 兼容 | 生态适配 |
| SDK | Python 单门 + 异步 | 协程友好 |
1.13 核心架构图(数据流视角)
这张图是 3.1 节的"逻辑版"——按"消息数据"从左到右贯穿整个栈,每一列代表一个层级。3.1 节强调调用栈的纵深,本节强调数据对象的横向流转。
关键看点:
- 入站对象(
InboundMessage)从左进入,经过 Channel Runtime 的适配,进入 MessageBus.inbound - AgentLoop 三件套(SessionLock / PendingQueue / AutoCompact + Checkpoint)协同保证会话一致性
- ContextBuilder 把 system prompt / history / 当前消息拼装成 LLM 可见的 messages
- AgentRunner 三件套(Hook chain / LLMProvider / ToolRegistry)按
for iteration循环反复调用 - TurnDelivery 把 LLM 流式输出 + 工具事件 + 思考段统一包装成
OutboundMessage.event - MessageBus.outbound 解耦到 ChannelManager,按事件类型分发
- 出站对象(
OutboundMessage)回到 Channel Runtime 的send()接口,发给聊天平台客户端
本文要点速查
- 架构核心:Channel → Bus → AgentLoop → AgentRunner → Provider 五层管道,每层职责单一可独立替换
- 关键设计:1s tick 兼顾巡检 + 6 阶段流水线(restore/compact/command/build/run/save/respond)+ SessionLock + ConcurrencyGate
- 实战建议:新接 Provider 只需实现
chat_stream()+count_tokens()两个方法即可接入完整 Agent 循环
tags:#nanobot #AI Agent #LLM #源码解析 #Python #消息总线 #异步 #事件驱动
更多推荐

所有评论(0)