PicoBot/docs/ARCHITECTURE.md
2026-07-14 13:02:06 +08:00

282 lines
13 KiB
Markdown
Raw Blame History

This file contains ambiguous Unicode characters

This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.

# PicoBot 架构
本文档描述 PicoBot 当前实现的运行时边界、数据流、并发模型和演进约束。它面向维护者和后续参与改进的 Agent是代码架构的主入口行为细节仍以代码和测试为最终依据。
## 1. 设计目标
PicoBot 是一个单进程、异步、可扩展的个人 AI 助手运行时。核心目标是:
- 用统一消息模型接入不同聊天渠道。
- 隔离渠道、会话、模型、工具和持久化职责。
- 同一会话内保持消息顺序,不同会话之间允许并发。
- 让外部 I/O、后台任务和程序关停都有明确边界。
- 通过 SQLite 保存会话、消息、记忆、定时任务及后台任务状态。
当前不是分布式系统。除调度任务使用数据库租约防止重复领取外,运行时会话状态由单个 Gateway 进程持有。
## 2. 运行模式与进程边界
PicoBot 只有一个二进制,提供两种模式:
| 模式 | 入口 | 职责 |
|------|------|------|
| Gateway | `cargo run -- gateway` | 组装服务、监听 HTTP/WebSocket、运行渠道、会话、调度器和后台任务 |
| CLI client | `cargo run -- chat` | 运行 Ratatui UI通过 WebSocket 使用 Gateway不持有业务状态 |
Gateway 启动时会切换进程工作目录到 `workspace_dir`。因此所有相对路径都应按 workspace 解释,不能假设仍位于源码仓库。
## 3. 组件关系
```mermaid
flowchart LR
External[CLI / Feishu] --> Channels[channels]
Channels -->|InboundMessage| Bus[MessageBus]
Bus --> Processor[Gateway message processor]
Processor --> Sessions[SessionManager]
Sessions --> Agent[AgentLoop]
Agent --> Providers[LLM providers]
Agent --> Tools[ToolRegistry / MCP]
Sessions <--> Storage[(SQLite)]
Sessions -->|OutboundMessage| Bus
Scheduler[Scheduler] --> Sessions
Bus --> Dispatcher[OutboundDispatcher]
Dispatcher --> Channels
Supervisor[TaskSupervisor] -. lifecycle .-> Processor
Supervisor -. lifecycle .-> Dispatcher
Supervisor -. lifecycle .-> Scheduler
Supervisor -. lifecycle .-> Sessions
```
### 模块职责
| 模块 | 拥有的职责 | 不应承担的职责 |
|------|------------|----------------|
| `gateway` | 依赖装配、HTTP/WS 入口、启动和关停顺序 | 业务规则、渠道协议细节 |
| `channels` | 外部协议适配、权限检查、媒体收发 | 会话选择、LLM 调用 |
| `bus` | 三条有界异步队列与出站投递协调 | 会话路由、业务状态 |
| `session` | dialog 路由、会话状态、串行工作队列、上下文和持久化协调 | 外部渠道协议 |
| `agent` | 单次无状态模型/工具循环、上下文压缩、子 Agent | 持有 dialog 生命周期 |
| `providers` | 把统一请求映射到模型 API | Session、Bus 或 Channel 感知 |
| `tools` / `mcp` | 工具定义、注册和执行适配 | 隐式修改会话路由 |
| `storage` | SQLite schema、迁移、原子 CRUD | 运行时调度策略 |
| `memory` | Knowledge/Timeline 的存取与召回 | 直接驱动消息发送 |
| `scheduler` | 领取到期任务、限并发执行、原子记录结果 | 复用聊天会话历史 |
| `task_supervisor` | 后台任务注册、取消、限时回收 | 业务级重试和结果语义 |
## 4. 消息与控制数据流
`MessageBus` 包含三条容量相同的 Tokio MPSC 队列:
- `inbound`Channel → Gateway message processor。
- `outbound`Session/Tool → `OutboundDispatcher`
- `control`WebSocket/Channel → Gateway message processor用于 dialog 操作。
### 普通消息
```mermaid
sequenceDiagram
participant C as Channel
participant B as MessageBus
participant G as Message processor
participant S as SessionManager
participant W as Per-session worker
participant A as AgentLoop
participant D as OutboundDispatcher
C->>B: publish InboundMessage
B->>G: consume inbound
G->>S: handle_message
S->>W: try_send AgentTask
S-->>G: AgentProcessing
W->>A: process(history)
A-->>W: final response
W->>B: publish OutboundMessage
B->>D: consume outbound
D->>C: Channel::send
```
关键语义:
- Gateway 的主消息处理循环不等待模型完成;普通消息进入对应 session worker 后立即返回 `AgentProcessing`
- 每个 session 有一条容量为 32 的队列,同一 session 串行处理,不同 session 的 worker 可并发执行。
- 队列满时明确拒绝新消息,不允许无界积压。
- Slash command 不进入 Agent 队列,由 `SessionManager` 直接执行,因此 `/stop` 等控制操作不会排在长模型调用之后。
### 出站投递
`OutboundDispatcher``(channel, chat_id)` 建立独立 lane
- 同一目标的消息保持顺序。
- 慢目标不会阻塞其他目标。
- 每条 lane 容量为 64空闲 300 秒后退出。
- 单次发送超时 30 秒;最多尝试 3 次,前两次失败后分别等待 1/2 秒。只有 `ConnectionError``SendError` 会重试。
- `deliver_outbound` 可等待渠道真实投递结果,等待上限 120 秒;普通 `publish_outbound` 只保证成功入队。
不要把“已进入 Bus”误认为“外部渠道已收到”。需要确认语义时必须使用 `deliver_outbound`
### Control 消息
WebSocket dialog 操作通过 `ControlMessage` 携带一次性回复通道。Gateway 在统一 message processor 中调用 `SessionManager`,再将 `SessionEvent` 回传给发起者。Bus 只承载消息,不解释操作。
## 5. 会话模型与并发不变量
Session ID 格式为:
```text
<channel>:<chat_id>:<dialog_id>
```
`SessionManagerInner` 保存:
- `sessions`:完整 Session ID 到内存 Session 的映射。
- `current_sessions``channel:chat_id` 到当前 dialog 的映射。
必须维护以下不变量:
1. 同一 Session 的 Agent 工作由一个 generation 对应的 worker 串行执行。
2. `/stop` 或 worker 替换会递增 `worker_generation`;旧 worker 不得再提交结果。
3. 慢操作模型、记忆召回、压缩、SQLite I/O不能长期持有 Session mutex。
4. 慢操作开始前记录 `state_version`,提交前重新验证,防止旧快照覆盖 `/clear``/delete` 等并发修改。
5. 持久化写入由 `persistence_lock` 串行化;多条相关记录应使用 Storage 的原子接口。
6. 内存先变更但持久化失败时,必须回滚精确匹配的消息后缀,不能删除无关的新状态。
SessionManager 负责组装会话上下文系统提示、Skills、召回的 Knowledge、压缩后的 Timeline 和当前消息历史。`AgentLoop` 接收完整输入执行一次模型/工具循环,本身不拥有会话状态。
## 6. 持久化
`Storage` 使用 SQLx + SQLite默认数据库为 `{workspace_dir}/picobot.db`。连接启用:
- WAL journal mode。
- foreign keys。
- 5 秒 busy timeout。
- schema version 迁移。
持久化范围包括 sessions、messages、memories、scheduled jobs、job runs 和 background tasks。修改 schema 时应:
1. 更新集中式 schema/迁移逻辑。
2. 保留已有数据库的升级路径。
3. 为新库初始化和旧库迁移分别增加测试。
4. 对“状态更新 + 执行记录”等复合写入使用事务。
### 安全边界
- API Key 和渠道凭据只来自配置占位符、`.env` 或进程环境,不得写入仓库。
- 日志不得输出 token、secret、Authorization header或包含临时凭据的完整 URL应记录脱敏后的 host/path 和必要诊断字段。
- Gateway 把 cwd 切到 workspace因此相对文件路径和 Shell 默认从 workspace 开始这不是硬沙箱。当前内置文件工具接受绝对路径Bash 也可访问进程权限允许的位置。若某场景需要硬边界,必须显式配置/实现 allowed directory 和进程隔离。
- `http_request``web_fetch` 的私网/回环地址校验属于 SSRF 防线,重构网络层时不能绕过。
- 外部内容、Tool 输出和 MCP 响应均是不可信输入;解析错误应返回结构化失败,不能 panic。
## 7. 后台任务与生命周期
`TaskSupervisor` 是 Gateway 内部后台任务的统一所有者。message processor、outbound dispatcher、scheduler、session workers、outbound lanes、通知消费者和子 Agent 后台任务都应通过它注册。
两种注册方式:
- `spawn`:收到全局取消后直接丢弃任务 future适合无需异步清理的任务。
- `spawn_graceful`:任务自己观察 cancellation token 并清理Supervisor 在总宽限期结束后再强制 abort。
新增长生命周期任务时必须满足:
- 有明确 owner禁止无法回收的裸 `tokio::spawn`
- 能响应取消;外部连接、重试 sleep 和阻塞式等待也要纳入取消分支。
- 等待任务退出必须有硬超时,超时后 abort 并回收 JoinHandle。
- 任务 panic、超时和永久错误要可观测且不能阻止其他组件清理。
WebSocket 每个客户端的 writer task 是连接局部任务,由连接 handler 自己限时回收;它不跨越连接生命周期。
## 8. 启动与关停顺序
### 启动
1. 加载配置和 `.env`,解析 workspace。
2. 创建并切换到 workspace。
3. 初始化 SQLite、MemoryManager、MessageBus 和 SessionManager。
4. 注册内置工具、渠道、MCP 工具和 Cron 工具。
5. 启动所有 Channel。
6. 通过 TaskSupervisor 启动 message processor、dispatcher 和 scheduler。
7. 绑定 Axum listener开始接收请求。
### 关停
1. `Ctrl-C` 触发 Axum graceful shutdown并取消所有 WebSocket 连接。
2. `ChannelManager::stop_all` 先停止外部消息入口并注销渠道。
3. 取消 TaskSupervisor停止接受新后台任务。
4. 在共享的 10 秒总宽限期内等待任务退出,之后 abort 剩余任务。
渠道自己的 `stop()` 也必须有界。以飞书为例端点请求、WebSocket 建连、重试等待和已连接循环共享 CancellationToken另有 5 秒强制回收兜底。
## 9. 扩展指南
### 新增 Channel
1. 实现 `Channel` trait仅处理外部协议和统一消息转换。
2.`ChannelManager::init` 注册,并通过同一个 MessageBus 收发。
3. 将可重试错误表示为 `ConnectionError`/`SendError`,永久错误使用其他类型。
4. 为 start/stop 幂等性、取消建连、投递失败和媒体边界增加测试。
5. 不要从 Channel 直接调用 SessionManager 或 Provider。
### 新增 Tool
1. 实现 `Tool`,在集中注册点加入 `ToolRegistry`
2. 参数 schema 必须明确;返回值保持可供模型消费的字符串协议。
3. 明确工具需要“workspace 默认目录”还是“不可逃逸的硬边界”;后者必须显式校验 canonical path不能只依赖 cwd。
4. 网络工具必须保留 SSRF/私网地址校验。
5. 长操作应有超时;后台执行应交给 SubAgentManager/TaskSupervisor。
### 新增 Provider
1. 实现 `LLMProvider`,保持其为纯 HTTP/API 适配器。
2. 统一映射文本、媒体、tool calls、usage 和错误。
3. 不在 Provider 中访问 Session、Bus 或 Channel。
4. 为请求序列化和响应兼容性编写离线单元测试。
### 修改 Session
先画出锁、慢 I/O、版本检查和持久化顺序。任何跨 `await` 的 Session 锁都需要特别审查;任何由旧快照产生的结果都必须在提交前验证 generation/state version。
## 10. 验证策略
按改动范围选择最小但充分的验证:
| 改动 | 至少执行 |
|------|----------|
| 文档 | 检查链接、命令和源码路径;`git diff --check` |
| Rust 实现 | 相关定向测试、`cargo test --lib``cargo clippy --all-targets --all-features -- -D warnings` |
| 构建/依赖 | 上述检查加 `cargo build` |
| Provider/真实渠道 | 离线测试;有凭据时再运行 ignored integration tests |
| SQLite schema | 新库测试、迁移测试、原子性测试 |
| 生命周期/并发 | 成功、失败、超时、取消、队列满和重复 start/stop 测试 |
模型 API 集成测试需要 `tests/test.env` 中的真实 API Key并默认 ignored`test_scheduler``test_request_format` 是可直接运行的离线测试。不应在无凭据时假装已经验证真实 Provider。
## 11. 演进决策清单
提交架构性修改前,逐项确认:
- 是否仍遵守 Channel → Bus → Session → Agent 的依赖方向?
- 是否引入第二个 owner、重复状态或绕过统一注册点
- 队列、并发、重试和等待是否全部有界?
- 取消信号能否覆盖建连、sleep、I/O 和清理阶段?
- 是否在持锁时执行了网络、模型或数据库慢操作?
- 内存与 SQLite 失败时能否保持一致?
- 错误是否区分瞬态和永久语义?
- 是否有测试覆盖正常路径与最危险的失败路径?
- README、AGENTS.md 和本文档是否需要同步更新?
## 12. 代码导航
| 主题 | 入口 |
|------|------|
| Gateway 装配和关停 | `src/gateway/mod.rs` |
| WebSocket 生命周期 | `src/gateway/ws.rs` |
| Bus 和消息类型 | `src/bus/mod.rs`, `src/bus/message.rs` |
| 出站并发与重试 | `src/bus/dispatcher.rs` |
| Channel 接口与注册 | `src/channels/base.rs`, `src/channels/manager.rs` |
| Session 核心 | `src/session/session.rs` |
| Session 命令/事件 | `src/session/commands.rs`, `src/session/events.rs` |
| Agent loop | `src/agent/agent_loop.rs` |
| 后台任务监督 | `src/task_supervisor.rs` |
| SQLite 初始化和迁移 | `src/storage/mod.rs` |
| Scheduler | `src/scheduler/mod.rs` |
| 配置加载 | `src/config/mod.rs` |