12 KiB
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 服务器,持有 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 发回原会话
关键约束
- 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
上下文压缩
当上下文接近 token 限制时触发:
- 快速裁剪:合并连续同角色消息,截断工具输出
- 硬截断:移除过老消息
- 压缩后保留用户消息确保结构完整
Skill 系统
三个优先级(高覆盖低):
{workspace}/skills/— 最高优先级~/.picobot/skills/— 中等优先级~/.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 IDcurrent_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 的处理原则:
- 短暂持 Session 锁抓取快照并记录
worker_generation/state_version。 - 释放锁后执行消息持久化、记忆召回、上下文压缩、LLM 和工具等慢操作。
- 提交由旧快照产生的结果前重新验证 generation/version,防止
/stop、/clear或/delete后写回陈旧状态。 - Session 持久化由独立
persistence_lock串行化;批量消息使用原子写入,失败时精确回滚内存后缀。 - 上下文溢出时按 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%) 时自动触发压缩:
- 快速裁剪:工具输出 ≥ 2000 字符时截断
- LLM 摘要:最多 3 轮,每轮找连续用户消息对,将中间的 assistant/tool 消息压缩为摘要 → 摘要作为 Timeline 记忆 持久化(importance 0.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:
- 按服务器配置连接
stdio、sse或streamable-http传输 - 调用 MCP
list_tools - 将每个 MCP tool 包装为
McpToolWrapper - 注册到当前 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 关停顺序:
- Ctrl-C 停止 Axum 接入并取消 WebSocket 连接。
ChannelManager::stop_all停止外部渠道并注销它们。- 取消 TaskSupervisor,并在共享的 10 秒宽限期内等待后台任务。
- 超时后 abort 剩余任务并回收 JoinHandle。
渠道 stop() 也必须有界。飞书端点请求、WebSocket 建连、重试 sleep 和连接循环共享 CancellationToken,另有 5 秒强制终止兜底。
当前斜杠命令
| 命令 | 说明 |
|---|---|
/new |
创建新对话 |
/sessions |
列出最近对话 |
/switch <dialog_id> |
切换到指定对话 |
/rename <title> |
重命名当前对话 |
/delete |
删除当前对话 |
/compact |
手动触发上下文压缩 |
/info |
显示当前对话信息 |
/dump |
保存当前对话为 markdown |
/?, /help |
显示帮助 |
/mcp |
显示 MCP 状态 |
/stop |
停止当前任务并清空消息队列 |