278 lines
13 KiB
Markdown
278 lines
13 KiB
Markdown
# PicoBot 架构机制
|
||
|
||
## 核心数据流
|
||
|
||
```
|
||
Channel → MessageBus.inbound → Gateway processor → SessionManager → per-session worker → AgentLoop
|
||
↑ │
|
||
└── Channel ← OutboundDispatcher ← per-conversation lane ← MessageBus.outbound
|
||
|
||
WebSocket/Channel → MessageBus.control → Gateway processor → SessionManager (dialog 操作)
|
||
Scheduler → SessionManager.handle_cron_message → AgentLoop → send_message
|
||
```
|
||
|
||
## 模块职责
|
||
|
||
| 模块 | 职责 |
|
||
|------|------|
|
||
| `gateway` | HTTP/WebSocket 服务器与嵌入式 WebUI,持有 GatewayState |
|
||
| `client` | TUI 聊天客户端 |
|
||
| `channels` | 外部集成(飞书、CLI),仅收发消息 |
|
||
| `bus` | 有界 inbound/outbound/control 队列;出站 dispatcher 与分目标 lane |
|
||
| `session` | 会话生命周期、dialog 操作、每 session 串行队列、上下文与持久化协调 |
|
||
| `agent` | LLM 调用循环、工具执行、上下文压缩、媒体处理、子 Agent |
|
||
| `providers` | LLM API 客户端(OpenAI 兼容、Anthropic) |
|
||
| `tools` | Agent 工具(bash、文件操作、搜索、HTTP、web、browser、memory、delegate 等) |
|
||
| `skills` | Skill 加载、管理和 prompt 构建 |
|
||
| `storage` | SQLite 持久化 |
|
||
| `scheduler` | Cron 作业调度 |
|
||
| `observability` | Observer 模式,agent/工具遥测事件 |
|
||
| `protocol` | WebSocket 协议消息定义 |
|
||
| `config` | 配置加载、环境变量替换、路径解析 |
|
||
| `memory` | 长期记忆存储与检索 |
|
||
| `mcp` | MCP(Model Context Protocol)工具集成 |
|
||
| `task_supervisor` | Gateway 后台任务注册、取消、限时等待和强制回收 |
|
||
|
||
## 功能边界
|
||
|
||
- Channels 仅收发消息,不感知 session 或 LLM
|
||
- MessageBus 本体持有三条有界队列;出站路由、顺序和重试由 `OutboundDispatcher` 负责
|
||
- SessionManager 拥有 session 状态、dialog 路由、上下文构建和每 session worker,并通过 worker 创建 AgentLoop
|
||
- AgentLoop 跨轮无状态,接收已准备的 history 调用 LLM、执行工具并返回一次结果
|
||
- Providers 是纯 HTTP 客户端,无 bus/session/channel 感知
|
||
- Tools 接收原始参数,返回字符串结果
|
||
- MCP 工具在 Gateway 初始化时连接服务器、发现工具,并包装成普通 Tool 注册到 ToolRegistry
|
||
- 子 Agent 由 `delegate` 工具创建,复用 provider 配置和按需过滤后的工具集;后台任务结果通过 MessageBus 发回原会话
|
||
- WebUI 聊天复用 `/ws` 与 `cli_chat`;同源管理 API 只读取受限的日志、任务、记忆,并对白名单配置文件做原子写入
|
||
- WebUI 使用 Svelte 5 + Vite,Bits UI 提供无样式可访问组件;`cargo build` 增量生成前端到 Cargo `OUT_DIR`,再嵌入单二进制,仓库不保存生成产物
|
||
|
||
## 关键约束
|
||
|
||
- Gateway 启动时切换到 workspace 目录
|
||
- SQLite 数据在 `{workspace}/picobot.db`
|
||
- ChannelManager 持有 MessageBus 和所有 channel
|
||
- OutboundDispatcher 通过 ChannelManager 路由出站消息
|
||
- Config `.env` 加载使用 `unsafe { env::set_var(...) }`
|
||
- `browser` 工具只有在 `browser.enabled=true` 时注册,依赖 Chrome/Chromium 与 WebDriver
|
||
- 同一 session 的普通消息串行处理,不同 session 可并发;session 队列容量为 32,满时明确拒绝
|
||
- 出站消息按 `(channel, chat_id)` 分 lane 保序;lane 容量为 64,慢目标不阻塞其他目标
|
||
- 长生命周期后台任务由 TaskSupervisor 管理;连接局部任务由其 owner 限时 join 或 abort
|
||
- 外部建连、重试等待和关停 join 必须可取消且有硬超时
|
||
- 不得记录 API Key、Authorization header 或包含临时凭据的完整连接 URL
|
||
- WebUI 当前无独立认证,默认回环监听是安全前提;对外暴露时必须由外层提供 TLS 和访问控制
|
||
|
||
## 上下文压缩
|
||
|
||
当上下文接近 token 限制时触发:
|
||
|
||
1. **快速裁剪**:合并连续同角色消息,截断工具输出
|
||
2. **硬截断**:移除过老消息
|
||
3. 压缩后保留用户消息确保结构完整
|
||
|
||
## Skill 系统
|
||
|
||
三个优先级(高覆盖低):
|
||
|
||
1. `{workspace}/skills/` — 最高优先级
|
||
2. `~/.picobot/skills/` — 中等优先级
|
||
3. `~/.agents/skills/` — 最低优先级
|
||
|
||
同名 skill 按优先级覆盖。每个 skill 是包含 `SKILL.md` 的目录。内置 skill 在 `~/.picobot/skills/` 下不存在时自动从二进制释放安装。
|
||
|
||
## 会话系统
|
||
|
||
### 会话 ID 格式
|
||
|
||
统一会话 ID 为三段式:**`<channel>:<chat_id>:<dialog_id>`**
|
||
|
||
| 部分 | 含义 | 示例 |
|
||
|------|------|------|
|
||
| `channel` | 消息渠道 | `cli_chat`、`feishu` |
|
||
| `chat_id` | 聊天/群组标识 | `sid_abc123` |
|
||
| `dialog_id` | 对话标识 | `default`、`d_xxxx`(短 ID) |
|
||
|
||
同一 `channel:chat_id` 下可有多个 dialog。`chat_scope()` 返回 `"channel:chat_id"` 用于分组。
|
||
|
||
### Session 生命周期
|
||
|
||
```
|
||
create → 存入 Storage → 载入 memory → 设为当前 dialog
|
||
↓
|
||
get_or_create
|
||
↓
|
||
← 接收消息、LLM 响应 →
|
||
↓
|
||
switch → rename → archive → delete(soft)
|
||
```
|
||
|
||
| 操作 | 效果 |
|
||
|------|------|
|
||
| `create` | 新建 dialog_id,立即持久化到 SQLite,设为当前 |
|
||
| `get_or_create` | 先在内存 HashMap 中找 → 再查 Storage → 都不存在则新建 |
|
||
| `switch_dialog` | 切换当前 dialog,目标 session 自动从 Storage 恢复入内存 |
|
||
| `list_dialogs` | 列出 `channel:chat_id` 下最近 10 个 session |
|
||
| `rename` | 更新标题,内存 + Storage 同步 |
|
||
| `delete` | 软删除(设 deleted_at),从内存移除 |
|
||
| `archive` | 设置 archived_at,从内存和当前 dialog 追踪中移除;可通过 include_archived 查询 |
|
||
|
||
### SessionManager 数据结构
|
||
|
||
两层追踪:
|
||
|
||
- **`sessions`**:`HashMap<String, Arc<Mutex<Session>>>` — 所有已加载的 session,key 为完整 session ID
|
||
- **`current_sessions`**:`HashMap<String, String>` — 每个 `channel:chat_id` 当前的 session ID
|
||
|
||
消息到达时 `resolve_dialog_id()` 按顺序确定接收 session:当前 session → Storage 最近活跃 session → 新建。
|
||
|
||
### 消息处理与并发
|
||
|
||
普通消息先 `try_send` 到该 session 的有界 worker 队列,Gateway 主 processor 随即返回 `AgentProcessing`。Slash command 直接执行,不进入此队列,因此 `/stop` 不会排在长模型调用后。
|
||
|
||
Worker 的处理原则:
|
||
|
||
1. 短暂持 Session 锁抓取快照并记录 `worker_generation`/`state_version`。
|
||
2. 释放锁后执行消息持久化、记忆召回、上下文压缩、LLM 和工具等慢操作。
|
||
3. 提交由旧快照产生的结果前重新验证 generation/version,防止 `/stop`、`/clear` 或 `/delete` 后写回陈旧状态。
|
||
4. Session 持久化由独立 `persistence_lock` 串行化;批量消息使用原子写入,失败时精确回滚内存后缀。
|
||
5. 上下文溢出时按 Provider 返回的真实限制重新压缩并重试。
|
||
|
||
### 会话恢复
|
||
|
||
从 Storage 恢复 session 时:
|
||
- 若 `last_compressed_message_at` 存在:先加载近 3 条 Timeline 记忆作为 `[Previous Context]`,再加载压缩标记后的原始消息
|
||
- 若无压缩记录:正常加载全部消息
|
||
- 自动修复断链的工具调用(gateway 崩溃中途重启导致)
|
||
|
||
---
|
||
|
||
## 记忆系统
|
||
|
||
### 记忆类别
|
||
|
||
| 类别 | 用途 | 生命周期 | 检索方式 |
|
||
|------|------|----------|----------|
|
||
| **Knowledge** | 事实、偏好、模式、洞察 | 长期保留,手动删除 | 每轮注入系统提示,关键词匹配 |
|
||
| **Timeline** | 历史会话摘要 | 配置预期保留 90 天;当前尚无自动清理循环 | `timeline_recall` 工具按需检索 |
|
||
|
||
### MemoryEntry 字段
|
||
|
||
| 字段 | 说明 |
|
||
|------|------|
|
||
| `id` | UUID |
|
||
| `key` | 唯一键,同 key 写入覆盖旧值 |
|
||
| `content` | 记忆内容 |
|
||
| `category` | `knowledge` 或 `timeline` |
|
||
| `importance` | 权重 (0.0–1.0),Timeline 默认为 0.3 |
|
||
| `session_id` | 关联会话(可选) |
|
||
|
||
### 存储与检索
|
||
|
||
- 主表 `memories` + FTS5 虚拟表 `memory_fts(key, content)` 全文索引
|
||
- 中文分词使用 jieba-rs 逐词精确匹配,用 OR 连接
|
||
- FTS5 无结果时回退到 LIKE 模糊匹配
|
||
- 支持 category、session_id、时间范围过滤
|
||
|
||
### 工作流程
|
||
|
||
```
|
||
用户消息到达
|
||
→ MemoryManager::recall(content, 5, Knowledge)
|
||
返回最多 5 条匹配的知识记忆(按 importance DESC)
|
||
→ 格式化为 "- key: content"
|
||
→ 作为运行时上下文附加到本轮 user message
|
||
→ LLM 可见,辅助回答
|
||
```
|
||
|
||
### 记忆工具
|
||
|
||
| 工具 | 写操作 | 说明 |
|
||
|------|--------|------|
|
||
| `memory_store` | 是 | 存储 Knowledge。必填: key, content。可选: importance |
|
||
| `memory_recall` | 否 | 搜索 Knowledge。必填: query(空格分隔关键词)。可选: since, until, limit |
|
||
| `timeline_recall` | 否 | 搜索 Timeline(压缩后的会话摘要)。必填: query。可选: session_id, since, until |
|
||
| `memory_forget` | 是 | 按 key 删除记忆 |
|
||
|
||
### 上下文压缩与 Timeline
|
||
|
||
LLM 对话上下文接近 token 限制 (默认 128K × 70%) 时自动触发压缩:
|
||
|
||
1. **快速裁剪**:工具输出 ≥ 2000 字符时截断
|
||
2. **LLM 摘要**:最多 3 轮,每轮找连续用户消息对,将中间的 assistant/tool 消息压缩为摘要 → 摘要作为 **Timeline 记忆** 持久化(importance 0.3)
|
||
3. **硬截断**:若仍超 90%,只保留前 N + 后 N 条消息
|
||
|
||
压缩后 `last_compressed_message_at` 标记边界,后续恢复时从标记点加载原始消息,以 Timeline 提供更早的上下文。
|
||
|
||
### 关键集成点
|
||
|
||
| 时机 | 操作 |
|
||
|------|------|
|
||
| 每次消息处理 | `memory_manager.recall()` 提取 Knowledge 上下文 |
|
||
| 系统提示构建 | `MemorySection` 渲染记忆工具指南;匹配的 Knowledge 附加到本轮 user message |
|
||
| 有压缩历史时 | `HistorySection` 提示 LLM 使用 `timeline_recall` |
|
||
| 压缩完成后 | 摘要自动存储为 Timeline 记忆 |
|
||
| 会话恢复 | 加载最近 Timeline 和压缩边界后的原始消息 |
|
||
|
||
`memory.recall_limit`、`idle_consolidation_minutes`、`timeline_retention_days` 和 `max_failures_before_degrade` 当前会被配置解析;其中每轮 Knowledge 召回在 worker 中仍固定为 5,其余自动维护策略尚未接入运行循环。不要把“配置可解析”误认为“行为已生效”。
|
||
|
||
---
|
||
|
||
## MCP 工具集成
|
||
|
||
Gateway 初始化时读取 `config.mcp.servers`:
|
||
|
||
1. 按服务器配置连接 `stdio`、`sse` 或 `streamable-http` 传输
|
||
2. 调用 MCP `list_tools`
|
||
3. 将每个 MCP tool 包装为 `McpToolWrapper`
|
||
4. 注册到当前 session 的 `ToolRegistry`
|
||
|
||
`/mcp` 斜杠命令会显示 MCP 服务器连接状态和工具列表。
|
||
|
||
---
|
||
|
||
## 子 Agent / delegate
|
||
|
||
`delegate` 工具用于把独立任务交给子 Agent:
|
||
|
||
| 模式 | 行为 |
|
||
|------|------|
|
||
| `inline` | 当前轮阻塞等待子 Agent 返回 |
|
||
| `background` | 后台运行,完成后通过原 channel/chat 通知 |
|
||
| `parallel` | 多个子 Agent 并发执行并聚合结果 |
|
||
|
||
默认工具集是只读工具:`file_read`、`file_search`、`content_search`、`web_fetch`、`http_request`、`calculator`。调用时可通过 `allowed_tools` 显式放开其他工具。后台任务会写入 `background_tasks` 表,默认 24 小时后清理。
|
||
|
||
后台子 Agent 通过 `TaskSupervisor::spawn_graceful` 注册,受 `gateway.max_concurrent_background_tasks` 限制;Gateway 关停时先收到取消信号,再在总宽限期内清理。
|
||
|
||
---
|
||
|
||
## 出站投递与关停
|
||
|
||
`OutboundDispatcher` 对每个 `(channel, chat_id)` 创建独立 lane:同一目标保持顺序,单次发送超时 30 秒,最多尝试 3 次。只有 `ConnectionError` 和 `SendError` 被视为瞬态错误;永久错误不重试。`deliver_outbound` 等待真实渠道投递结果,`publish_outbound` 只表示成功入队。
|
||
|
||
Gateway 关停顺序:
|
||
|
||
1. Ctrl-C 停止 Axum 接入并取消 WebSocket 连接。
|
||
2. `ChannelManager::stop_all` 停止外部渠道并注销它们。
|
||
3. 取消 TaskSupervisor,并在共享的 10 秒宽限期内等待后台任务。
|
||
4. 超时后 abort 剩余任务并回收 JoinHandle。
|
||
|
||
渠道 `stop()` 也必须有界。飞书端点请求、WebSocket 建连、重试 sleep 和连接循环共享 CancellationToken,另有 5 秒强制终止兜底。
|
||
|
||
---
|
||
|
||
## 当前斜杠命令
|
||
|
||
| 命令 | 说明 |
|
||
|------|------|
|
||
| `/new` | 创建新对话 |
|
||
| `/sessions` | 列出最近对话 |
|
||
| `/switch <dialog_id>` | 切换到指定对话 |
|
||
| `/rename <title>` | 重命名当前对话 |
|
||
| `/delete` | 删除当前对话 |
|
||
| `/compact` | 手动触发上下文压缩 |
|
||
| `/info` | 显示当前对话信息 |
|
||
| `/dump` | 保存当前对话为 markdown |
|
||
| `/?`, `/help` | 显示帮助 |
|
||
| `/mcp` | 显示 MCP 状态 |
|
||
| `/stop` | 停止当前任务并清空消息队列 |
|