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 给出响应"的端到端架构图:

通道出口

Hook与投递

Provider

AgentRunner

AgentLoop

总线层

通道层

客户端

/sendMessage

publish_inbound(InboundMessage)

consume_inbound()

yield chunks

tool_calls

ToolResult

注入

继承

包裹

publish_outbound

consume_outbound

channel.send()

Telegram 客户端

TelegramChannel
runtime.py
鉴权 / 解析 / 包装

MessageBus.inbound
queue.py
asyncio.Queue

AgentLoop.run()
1s tick 等待消息

_dispatch
SessionLock + ConcurrencyGate

_process_message
六阶段流水线

1. restore
checkpoint / pending_user_turn

2. compact
AutoCompact.prepare_session

3. command
/stop /new 短路

4. build
ContextBuilder

5. run
AgentRunner.run()

6. save
_persist_turn → SessionManager

for iteration in range(max)

_request_model
chat_stream_with_retry

_execute_tools
按 concurrency_safe 切批

_try_drain_injections

LLMProvider 抽象

AnthropicProvider
AsyncAnthropic SDK

OpenAICompatProvider
AsyncOpenAI SDK

FallbackProvider
circuit breaker

CompositeHook
AgentProgressHook + FileEditActivityHook

TurnDelivery
route + stream_id + segment

MessageBus.outbound
queue.py

ChannelManager._dispatch_outbound
按 event 类型分发
合并 / 抑制 / retry

关键路径解读

  1. 入站路径(蓝色流程):Telegram 客户端 → TelegramChannel 适配器 → MessageBus.inbound → AgentLoop.run() 1s tick
  2. 编排路径:AgentLoop → _dispatch(SessionLock 串行同 session / ConcurrencyGate 限制跨 session 并发) → _process_message 六阶段流水线
  3. 执行路径:AgentRunner 的 for iteration 循环 → _request_model 调 LLM → 流式 chunks 通过 CompositeHook 翻译成 OutboundMessage → 工具调用走 _execute_tools → 注入消费走 _try_drain_injections
  4. Provider 路径LLMProvider 抽象 → AnthropicProvider / OpenAICompatProvider / FallbackProvider(装饰器模式)
  5. 出站路径TurnDelivery 投递 → MessageBus.outbound → ChannelManager._dispatch_outbound 按事件类型分发(合并 / 抑制 / retry) → channel.send() → 回到 Telegram 客户端

1.2 关键路径详解

1.2.1 入站路径:Channel → AgentLoop

AgentLoop MessageBus PairingStore TelegramChannel Telegram 客户端 AgentLoop MessageBus PairingStore TelegramChannel Telegram 客户端 alt [未授权 + DM] queue.put / 非阻塞 update (sendMessage) 1 is_allowed(sender) 权限检查 2 is_approved(channel, sender_id) 3 False 4 generate_code() 一次性码 5 提示配对码 6 publish_inbound(InboundMessage) 7 consume_inbound() 8 InboundMessage 9

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

事件类型

LLM响应

on_content_delta

on_thinking_delta

on_tool_call_delta

on_tool_call_done

on_stream_end

emit_reasoning_end

StreamDeltaEvent

ProgressEvent
reasoning_delta

ProgressEvent
tool_events start

ProgressEvent
tool_events end

StreamEndEvent

ProgressEvent
reasoning_end

核心要点速查(建议收藏)

  • 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 上下文注入机制

AgentLoopcontextvars 把每个 turn 的运行时状态注入到任意深度的调用栈:

enter

enter

enter

任意深处

任意深处

任意深处

AgentLoop._run_turn

bind_request_context
channel/chat_id/session_key/...
(ContextVar)

bind_workspace_scope
WorkspaceScope
(ContextVar)

bind_file_states
FileStates
(ContextVar)

_run_agent_loop

current_request_context()

current_workspace_scope()

current_file_states()

四个关键 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"

崩溃后启动

正常路径

用户消息已写但 assistant 未生成

启动

有 _pending_user_turn?

插入占位 assistant 消息
'Previous assistant message omitted'

清 _pending_user_turn 标记

用户消息

_persist_user_message

LLM 调用

有 tool_call?

执行工具

_set_runtime_checkpoint

终止

启动 / 收到新消息

_restore_turn

有 checkpoint?

回填已完成的 tool_result

重放 pending_tool_calls

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 模式:

写状态

创建 .tmp 临时文件

f.write(content)

f.flush()

os.fsync(f.fileno())
强制刷盘

os.replace(tmp, final_path)
原子替换

父目录 fsync

写入完成

# 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 抽风而无限循环:

LLM 抽风风险

层 1: max_iterations
默认 200
AgentLoop 注入 AgentRunSpec

层 2: _MAX_INJECTION_CYCLES = 5
runner.py 限制单回合注入

层 3: _MAX_GOAL_CONTINUATION_ROUNDS = 12
turn_continuation.py 限制持续目标

强制退出

  1. max_iterations(默认 200)由 AgentLoop 注入到 AgentRunSpec.max_iterationsAgentRunner 跑到上限就强制退出
  2. _MAX_INJECTION_CYCLES = 5 在 [nanobot/agent/runner.py]限制单回合最多 5 次注入
  3. _MAX_GOAL_CONTINUATION_ROUNDS = 12 在 [nanobot/session/turn_continuation.py] 限制持续目标最多 12 轮续传

1.7 并发模型

单 Turn 层

单 Session 层

全局层

限流

ConcurrencyGate
asyncio.Semaphore
默认 max=3

所有 session
共享一个信号量

_session_locks
dict[str, asyncio.Lock]

同 session_key
串行进入 critical section

_pending_queues
asyncio.Queue maxsize=20

mid-turn 注入
不抢占,走队列

三层并发语义

  • 同一 session_key 串行(_session_locks
  • 跨 sessionConcurrencyGate(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

回调链路:

awatch 检测

config.json 改写

on_change 回调

agent.invalidate_runtime_config()

清空内部缓存

下次 resolve_runtime_config()

重新加载

全链路无需重启

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

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()

清理顺序:

SIGINT / SIGTERM
二次按下

cron.stop()

agent.stop()

channels.stop_all()

取消所有 _background_tasks

agent.sessions.flush_all()

关键: 把缓存会话 fsync 到磁盘

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;OpenAICompatProviderAsyncOpenAI SDK 处理 30+ provider。

好处:统一重试 / fallback / 流式语义。
代价:必须满足"相同的最小字段"——但现实是各家 schema 差异大,导致 OpenAICompatProvider 1750 行。

权衡 3:Tool 自发现 vs 显式注册

ToolLoaderpkgutil.iter_modules 扫描 nanobot/agent/tools/ 下所有模块。但工具作者必须继承 Tool 并实现 execute

好处:新工具"装完即用"。
代价:测试时容易意外加载所有工具——所以 _SKIP_MODULES 必须明确列出。

权衡 4:持久化选 JSONL + git

MemoryStoreMEMORY.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 节强调调用栈的纵深,本节强调数据对象的横向流转。

TurnDelivery

Tool 抽象

AgentRunner
纯 LLM 循环

AgentLoop 编排核心

MessageBus
32 行的极简双队列

consume

chunks

response

execute

ToolResult

publish_outbound

publish_inbound

AgentRunSpec

consume

channel.send()

ChannelManager._dispatch_outbound

按 event 类型分发
合并 / 抑制 / retry

OutboundMessage (发往通道)

event:
Progress / StreamDelta / StreamEnd / StreamedResponse

content + metadata

ContextBuilder

build_system_prompt
8 段拼接

build_messages

Channel Runtime

Telegram / Discord / Slack / ...
17+ 自发现通道

InboundMessage (从通道进入)

channel / sender_id / chat_id

content / media

metadata / session_key_override

inbound: asyncio.Queue

outbound: asyncio.Queue

SessionLock
+ PendingQueue

_process_message
6 阶段流水线

AutoCompact
+ Checkpoint

Hook chain
(Composite)

LLMProvider
(Anthropic / OpenAI / Fallback)

ToolRegistry

filesystem / shell / web / MCP / ...
17+ 自发现工具

route + stream_id + segment

关键看点

  • 入站对象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() 接口,发给聊天平台客户端

本文要点速查

  1. 架构核心:Channel → Bus → AgentLoop → AgentRunner → Provider 五层管道,每层职责单一可独立替换
  2. 关键设计:1s tick 兼顾巡检 + 6 阶段流水线(restore/compact/command/build/run/save/respond)+ SessionLock + ConcurrencyGate
  3. 实战建议:新接 Provider 只需实现 chat_stream() + count_tokens() 两个方法即可接入完整 Agent 循环

tags:#nanobot #AI Agent #LLM #源码解析 #Python #消息总线 #异步 #事件驱动

Logo

Agent 垂直技术社区,欢迎活跃、内容共建。

更多推荐