300 lines
16 KiB
Markdown
Raw Permalink 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 架构机制
## 核心数据流
```
Channel → MessageBus.inbound → Gateway processor → SessionManager → per-session worker → AgentLoop
↑ │
└── Channel ← OutboundDispatcher ← per-conversation lane ← MessageBus.outbound
AgentLoop → TurnEvent → TurnController → latest TurnSnapshot → DeliveryCoordinator → TurnSink → Channel
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 串行队列、Turn 状态、上下文与持久化协调 |
| `agent` | LLM 调用循环、工具执行、上下文压缩、媒体处理、子 Agent、Turn 语义事件 |
| `providers` | OpenAI/Anthropic 原生流解析统一正文、reasoning、工具、usage 与私有回放状态 |
| `delivery` | 完整 Turn 快照的展示过滤、latest-wins 节流、终态投递和 TurnSink 生命周期 |
| `tools` | Agent 工具bash、文件操作、搜索、HTTP、web、browser、memory、delegate 等) |
| `skills` | Skill 加载、管理和 prompt 构建 |
| `storage` | SQLite 持久化 |
| `scheduler` | Cron 作业调度 |
| `observability` | Observer 模式agent/工具遥测事件 |
| `protocol` | WebSocket 协议消息定义 |
| `config` | 配置加载、环境变量替换、路径解析 |
| `memory` | 长期记忆存储与检索 |
| `mcp` | MCPModel Context Protocol工具集成 |
| `task_supervisor` | Gateway 后台任务注册、取消、限时等待和强制回收 |
| `work` | 每个 session 的单 active plan、并行子项状态、版本与变更事件 |
## 功能边界
- Channels 通过 MessageBus 发布入站消息,通过 OutboundDispatcher 或每 Turn 一个的 TurnSink 接收出站写入,不感知 session 或 LLM
- MessageBus 本体持有三条有界队列;出站路由、顺序和重试由 `OutboundDispatcher` 负责
- SessionManager 拥有 session 状态、dialog 路由、上下文构建和每 session worker并通过 worker 创建 AgentLoop
- TurnController 是活动 Turn 状态的唯一 ownerSession 在消息原子提交成功后才发布 Completed
- AgentLoop 跨轮无状态,接收已准备的 history 调用 LLM、执行工具并返回一次结果
- Providers 是纯 HTTP 流客户端,无 bus/session/channel 感知;签名 reasoning 状态只回放给匹配 Provider不下发客户端或 Channel
- DeliveryCoordinator 只投影完整快照,不修改会话历史;慢消费者跳过中间 revision终态显式、有界投递
- 每个活动 Turn 独占一个 TurnSink平台 message ID 和 reaction 清理状态只存在于 sink 内
- Tools 接收原始参数,返回字符串结果
- MCP 工具在 Gateway 初始化时连接服务器、发现工具,并包装成普通 Tool 注册到 ToolRegistry
- 子 Agent 由 `delegate` 工具创建,复用 provider 配置和按需过滤后的工具集;后台任务结果通过 MessageBus 发回原会话
- 复杂任务可使用 `todo` 创建 session 级计划;多个子项可通过 `delegate.plan_item_id` 并行委托,子 Agent 不能修改计划
- WebUI 聊天复用 `/ws``cli_chat`;同源管理 API 只读取受限的日志、任务、记忆,并对白名单配置文件做原子写入
- WebUI 使用 Svelte 5 + ViteBits UI 提供无样式可访问组件;`cargo build` 增量生成前端到 Cargo `OUT_DIR`,再嵌入单二进制,仓库不保存生成产物
## 关键约束
- Gateway 启动时切换到 workspace 目录
- SQLite 数据在 `{workspace}/picobot.db`
- ChannelManager 持有 MessageBus 和所有 channel
- OutboundDispatcher 通过 ChannelManager 路由出站消息
- 配置目录 `.env` 与 workspace `.env` 仅在单线程启动阶段分层加载,并使用 `unsafe { env::set_var(...) }` 写入进程环境;优先级为既有进程环境 > workspace > 配置目录
- `browser` 工具只有在 `browser.enabled=true` 时注册,依赖 Chrome/Chromium 与 WebDriver
- 同一 session 的普通消息串行处理,不同 session 可并发session 队列容量为 32满时明确拒绝
- 出站消息按 `(channel, chat_id)` 分 lane 保序lane 容量为 64慢目标不阻塞其他目标
- 活动 Turn 与普通出站消息共享 `(channel, chat_id)` 写锁;禁止把 token delta 放入 MessageBus
- `cli_chat` 向 TUI/WebUI 发送统一 `turn_updated` 完整快照;飞书默认 FinalOnly开启 `live_updates` 后编辑同一卡片
- 长生命周期后台任务由 TaskSupervisor 管理;连接局部任务由其 owner 限时 join 或 abort
- 外部建连、重试等待和关停 join 必须可取消且有硬超时
- 不得记录 API Key、Authorization header 或包含临时凭据的完整连接 URL
- WebUI 管理 API 与 `/ws` 默认要求设备配对;一次性代码由本机 CLI 签发,服务端只持久化令牌哈希。对外暴露时仍必须由外层提供 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>>>` — 所有已加载的 sessionkey 为完整 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 返回的真实限制重新压缩并重试。
WebUI/TUI 的 Active Turn 使用 `send_message(files=...)` 向自身 session 投递附件时,附件暂存到 task-local Turn delivery成功结束后并入最终 assistant 消息,因此工具链始终排在附件回复之前且不会出现自引用来源前缀。其他自投递要求 task-local Turn ID 与 session 的 active Turn 匹配;历史中的 assistant/system 附件只作为文本清单提供给模型,原生媒体块仅用于 user 输入和当前工具结果。
### 活动 Turn
每个主 Agent 请求会创建一个内存 Turn。Provider delta 经 AgentLoop 转换为 reasoning、正文、工具开始/完成等语义事件TurnController 归约为有序 block 和单调 revision 的完整快照。TUI/WebUI 使用 `history + active_turn` 渲染,不自行拼接 token中间帧可丢下一快照会自动收敛。
展示策略在 Gateway 核心出口应用:交互客户端可显示 reasoning 和详细工具状态;外部渠道隐藏 reasoning、工具仅显示紧凑状态无人值守投递只保留正文。运行态不逐 token 入库,完成、取消或中断时才原子保存消息及 completion status。
### 会话恢复
从 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.01.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 关停时先收到取消信号,再在总宽限期内清理。
## Session Todo 计划
每个 session 最多有一个 active plan普通闲聊没有计划上下文。计划和子项分别持久化到 `task_plans``task_items`历史压缩后仍从权威状态生成精简摘要。WebUI 聊天页通过 `session_plan``plan_updated` WebSocket 帧显示默认隐藏的侧栏,当前 session 变化时自动展开,其他 session 只标记未读。
---
## 出站投递与关停
`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` | 停止当前任务并清空消息队列 |