docs+test: 补充架构文档与 anthropic provider 单测
P3 文档: - 新增 ARCHITECTURE.md,聚焦数据流与 7 个关键设计决策 (MessageBus 解耦/SessionPool 隔离/AgentLoop 循环/SQLite 池化/重启机制/嵌入静态文件/Safety Guard) - 不写代码导读,只写'为什么这样设计',代码是唯一真相源 P2 测试: - src/providers/anthropic.rs 曾零测试(396 行),补 15 个纯函数单测 - 覆盖:data URL 解析、图片过滤逻辑、内部字段过滤、响应反序列化、错误链格式化 - providers 模块测试密度 31% -> 提升,重点补齐 anthropic 空白
This commit is contained in:
parent
c724bbf864
commit
e6b2fcfb6d
116
ARCHITECTURE.md
Normal file
116
ARCHITECTURE.md
Normal file
@ -0,0 +1,116 @@
|
|||||||
|
# PicoBot 架构
|
||||||
|
|
||||||
|
> 本文档聚焦"为什么这样设计"和"数据如何流动",不是代码导读。
|
||||||
|
> 代码是唯一真相源,文档可能滞后;有冲突以代码为准。
|
||||||
|
|
||||||
|
## 一句话定位
|
||||||
|
|
||||||
|
PicoBot 是一个**多渠道接入的 Agent 网关**:外部消息(微信/飞书/Web/CLI)经统一总线进入,
|
||||||
|
由会话管理器路由到对应 Agent 循环,Agent 调用 LLM + 工具完成任务,结果原路返回。
|
||||||
|
|
||||||
|
## 核心数据流
|
||||||
|
|
||||||
|
```
|
||||||
|
┌─────────┐ ┌──────────┐ ┌──────────────┐ ┌────────────┐ ┌──────────┐
|
||||||
|
│ Channel │──▶│ MessageBus│──▶│InboundProc. │──▶│ Session │──▶│AgentLoop │
|
||||||
|
│ (微信/ │ │ (解耦) │ │(并发+限流) │ │ Manager │ │(LLM+工具)│
|
||||||
|
│ 飞书/Web)│ │ │ │ │ │ │ │ │
|
||||||
|
└─────────┘ └──────────┘ └──────────────┘ └────────────┘ └────┬─────┘
|
||||||
|
▲ │
|
||||||
|
│ ▼
|
||||||
|
┌────┴────┐ ┌──────────┐ ┌────────────┐ ┌──────────┐
|
||||||
|
│Outbound │◀──│ MessageBus│◀─── SessionMessageSender ◀─│ storage │◀─│ tools │
|
||||||
|
│Dispatcher│ │ │ │ (SQLite) │ │(bash/file│
|
||||||
|
└─────────┘ └──────────┘ └────────────┘ │ /memory) │
|
||||||
|
└──────────┘
|
||||||
|
```
|
||||||
|
|
||||||
|
**入站**:Channel → MessageBus → InboundProcessor(并发控制 + 限流)→ SessionManager → AgentLoop
|
||||||
|
**出站**:AgentLoop → SessionMessageSender → MessageBus → OutboundDispatcher → Channel
|
||||||
|
**持久化**:AgentLoop / Tools → SessionStore(SQLite + r2d2 连接池)
|
||||||
|
|
||||||
|
## 关键设计决策
|
||||||
|
|
||||||
|
### 1. MessageBus 解耦 Channel 与 Session
|
||||||
|
|
||||||
|
**问题**:多个 Channel(微信/飞书/Web)接入,每个 Channel 协议不同,但 Session 处理逻辑相同。
|
||||||
|
**决策**:引入 MessageBus 作为中间件,Channel 只负责协议适配和收发,Session 不关心消息来自哪个 Channel。
|
||||||
|
**代价**:多一层间接。换来的是新增 Channel(如钉钉)只需实现 Channel trait,不碰 Session 逻辑。
|
||||||
|
|
||||||
|
### 2. SessionPool 会话隔离
|
||||||
|
|
||||||
|
**问题**:多用户同时对话,会话状态不能串扰。
|
||||||
|
**决策**:SessionManager 持有 SessionPool,按 (channel_name, chat_id) 路由到独立 Session。
|
||||||
|
每个 Session 有自己的 AgentLoop、消息历史、工具上下文。
|
||||||
|
**代价**:内存占用随活跃会话数增长。用 session_ttl_hours 过期回收。
|
||||||
|
|
||||||
|
### 3. AgentLoop 的工具循环
|
||||||
|
|
||||||
|
**问题**:LLM 需要多轮工具调用才能完成任务(如"读文件→分析→写文件")。
|
||||||
|
**决策**:AgentLoop 是一个有界循环(max_tool_iterations,默认 1000),每轮:
|
||||||
|
1. 把消息历史 + 工具定义发给 LLM
|
||||||
|
2. LLM 返回文本或 tool_call
|
||||||
|
3. 如果是 tool_call,执行工具,把结果加入历史,回到 1
|
||||||
|
4. 如果是文本,结束循环
|
||||||
|
**代价**:单次对话可能很长。用 CancelManager 支持中途取消。
|
||||||
|
|
||||||
|
### 4. SQLite + r2d2 连接池
|
||||||
|
|
||||||
|
**问题**:需要持久化会话历史、话题、记忆、待办,且要支持并发读写的 sub-agent 场景。
|
||||||
|
**决策**:SQLite(bundled)+ r2d2 连接池(max_size=8)+ busy_timeout(30s)。
|
||||||
|
**代价**:SQLite 写并发有限。通过 busy_timeout + 事务隔离级别(之前的修复)缓解锁冲突。
|
||||||
|
**不选 Postgres 的原因**:单机部署、零外部依赖、足够用。
|
||||||
|
|
||||||
|
### 5. 重启机制(watch channel)
|
||||||
|
|
||||||
|
**问题**:配置变更后需要重启 gateway,但不能要求用户手动杀进程。
|
||||||
|
**决策**:main.rs 用 `while should_restart` 循环,gateway 通过 `watch::Sender<bool>` 通知是否需要重启。
|
||||||
|
配置 API `/api/restart` 触发 graceful shutdown,main 收到 should_restart=true 后重新初始化。
|
||||||
|
**代价**:重启期间短暂不可用。比热重载简单且可靠。
|
||||||
|
|
||||||
|
### 6. 嵌入式静态文件
|
||||||
|
|
||||||
|
**问题**:Web 前端需要随二进制分发,但不想要求用户额外下载。
|
||||||
|
**决策**:build.rs 在 cargo build 时执行 npm run build,产物通过 rust-embed 编译进二进制。
|
||||||
|
开发时设 `STATIC_DIR` 环境变量走磁盘文件,支持热更新。
|
||||||
|
**代价**:二进制体积增大。换来的是单文件部署。
|
||||||
|
|
||||||
|
### 7. Safety Guard(命令安全护栏)
|
||||||
|
|
||||||
|
**问题**:Agent 可以调用 bash 工具执行任意命令,需要防止误操作(如 `format C:`、`rm -rf`)。
|
||||||
|
**决策**:platform 模块按平台注入危险命令正则,执行前匹配拦截。
|
||||||
|
**关键教训**:正则要精确(曾因 `\bformat\s+` 误拦 `dart format`),按平台分组避免跨平台误伤。
|
||||||
|
|
||||||
|
## 模块职责速查
|
||||||
|
|
||||||
|
| 模块 | 职责 | 关键文件 |
|
||||||
|
|------|------|---------|
|
||||||
|
| gateway | HTTP/WS 服务、路由、生命周期 | `gateway/mod.rs`, `gateway/runtime.rs` |
|
||||||
|
| gateway/session | 会话管理、路由、池化 | `gateway/session.rs`, `session_pool.rs` |
|
||||||
|
| gateway/processor | 入站消息处理、并发控制 | `gateway/processor.rs` |
|
||||||
|
| agent | Agent 循环、上下文压缩 | `agent/agent_loop.rs` |
|
||||||
|
| providers | LLM Provider 抽象(OpenAI/Anthropic) | `providers/openai.rs`, `anthropic.rs` |
|
||||||
|
| tools | 工具实现与注册 | `tools/` (bash/file/memory/task/...) |
|
||||||
|
| storage | SQLite 持久化 | `storage/mod.rs`, `migrations.rs` |
|
||||||
|
| channels | 渠道适配(微信/飞书/CLI) | `channels/` |
|
||||||
|
| bus | 消息总线(解耦 channel 与 session) | `bus/message.rs` |
|
||||||
|
| command | 前端命令处理(话题/会话/记忆 CRUD) | `command/handlers/` |
|
||||||
|
| mcp | Model Context Protocol 客户端 | `mcp/client.rs` |
|
||||||
|
| scheduler | 定时任务调度 | `scheduler/mod.rs` |
|
||||||
|
| skills | 技能加载与激活 | `skills/mod.rs` |
|
||||||
|
| experts | 专家配置运行时 | `experts/mod.rs` |
|
||||||
|
|
||||||
|
## 测试策略
|
||||||
|
|
||||||
|
- **单元测试**:与代码同文件 `#[cfg(test)] mod tests`,覆盖纯函数和逻辑分支
|
||||||
|
- **集成测试**:`tests/` 目录,覆盖跨模块请求格式
|
||||||
|
- **测试密度**(test 行 / 总行):tools 98.8%、gateway 96.4%、storage 91.8% 为高覆盖区;
|
||||||
|
providers 31%(anthropic 曾为 0%,已补)、bus 18.7%、cli 4.5% 为薄弱区
|
||||||
|
- **不强制覆盖率工具**:静态审计 + 高风险区定向补测,比全量 tarpaulin 更务实
|
||||||
|
|
||||||
|
## 工程化基线
|
||||||
|
|
||||||
|
- **格式化**:rustfmt(Rust)+ prettier(前端),CI 强制 `--check`
|
||||||
|
- **静态检查**:clippy(Rust)+ eslint(前端),CI 强制
|
||||||
|
- **CI**:GitHub Actions,双平台(ubuntu + windows)跑 fmt + clippy + test + eslint + tsc + vitest
|
||||||
|
- **本地**:`make check` 与 CI 完全对齐
|
||||||
@ -400,3 +400,229 @@ impl LLMProvider for AnthropicProvider {
|
|||||||
&self.model_id
|
&self.model_id
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
#[cfg(test)]
|
||||||
|
mod tests {
|
||||||
|
use super::*;
|
||||||
|
use crate::domain::messages::ContentBlock;
|
||||||
|
use std::collections::HashMap;
|
||||||
|
|
||||||
|
/// 构造一个最小 Provider 用于测试配置驱动的方法
|
||||||
|
fn make_provider(model_extra: HashMap<String, serde_json::Value>) -> AnthropicProvider {
|
||||||
|
AnthropicProvider::new(
|
||||||
|
"test".to_string(),
|
||||||
|
"key".to_string(),
|
||||||
|
"https://api.test".to_string(),
|
||||||
|
HashMap::new(),
|
||||||
|
30,
|
||||||
|
"claude-test".to_string(),
|
||||||
|
None,
|
||||||
|
None,
|
||||||
|
model_extra,
|
||||||
|
)
|
||||||
|
}
|
||||||
|
|
||||||
|
// ---- convert_image_url_to_anthropic ----
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn test_convert_data_url_extracts_media_type_and_base64() {
|
||||||
|
let url = "data:image/png;base64,iVBORw0KGgo=";
|
||||||
|
let v = convert_image_url_to_anthropic(url);
|
||||||
|
assert_eq!(v["type"], "image");
|
||||||
|
assert_eq!(v["source"]["type"], "base64");
|
||||||
|
assert_eq!(v["source"]["media_type"], "image/png");
|
||||||
|
assert_eq!(v["source"]["data"], "iVBORw0KGgo=");
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn test_convert_data_url_jpeg() {
|
||||||
|
let url = "data:image/jpeg;base64,/9j/4AAQ";
|
||||||
|
let v = convert_image_url_to_anthropic(url);
|
||||||
|
assert_eq!(v["source"]["media_type"], "image/jpeg");
|
||||||
|
assert_eq!(v["source"]["data"], "/9j/4AAQ");
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn test_convert_regular_url_uses_url_source() {
|
||||||
|
let url = "https://example.com/img.png";
|
||||||
|
let v = convert_image_url_to_anthropic(url);
|
||||||
|
assert_eq!(v["type"], "image");
|
||||||
|
assert_eq!(v["source"]["type"], "url");
|
||||||
|
assert_eq!(v["source"]["url"], url);
|
||||||
|
}
|
||||||
|
|
||||||
|
// ---- convert_content_blocks: 图片不支持时的过滤 ----
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn test_convert_blocks_filters_images_when_unsupported() {
|
||||||
|
let blocks = vec![
|
||||||
|
ContentBlock::text("hello"),
|
||||||
|
ContentBlock::image_url("data:image/png;base64,abc"),
|
||||||
|
ContentBlock::image_url("data:image/png;base64,def"),
|
||||||
|
];
|
||||||
|
let result = convert_content_blocks(false, "test", "claude-test", &blocks, 0);
|
||||||
|
// 文本块保留,图片块被替换为通知
|
||||||
|
assert_eq!(result.len(), 2);
|
||||||
|
assert_eq!(result[0]["type"], "text");
|
||||||
|
assert_eq!(result[0]["text"], "hello");
|
||||||
|
// 第二个是合并的图片通知
|
||||||
|
assert_eq!(result[1]["type"], "text");
|
||||||
|
let notice = result[1]["text"].as_str().unwrap();
|
||||||
|
assert!(notice.contains("第 1 张图片"));
|
||||||
|
assert!(notice.contains("第 2 张图片"));
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn test_convert_blocks_keeps_images_when_supported() {
|
||||||
|
let blocks = vec![
|
||||||
|
ContentBlock::text("hi"),
|
||||||
|
ContentBlock::image_url("data:image/png;base64,abc"),
|
||||||
|
];
|
||||||
|
let result = convert_content_blocks(true, "test", "claude-test", &blocks, 0);
|
||||||
|
assert_eq!(result.len(), 2);
|
||||||
|
assert_eq!(result[0]["type"], "text");
|
||||||
|
assert_eq!(result[1]["type"], "image");
|
||||||
|
assert_eq!(result[1]["source"]["data"], "abc");
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn test_convert_blocks_text_only_passthrough() {
|
||||||
|
let blocks = vec![ContentBlock::text("just text")];
|
||||||
|
let result = convert_content_blocks(false, "test", "claude-test", &blocks, 0);
|
||||||
|
assert_eq!(result.len(), 1);
|
||||||
|
assert_eq!(result[0]["type"], "text");
|
||||||
|
}
|
||||||
|
|
||||||
|
// ---- request_model_extra: 内部字段过滤 ----
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn test_request_model_extra_filters_internal_keys() {
|
||||||
|
let mut extra = HashMap::new();
|
||||||
|
extra.insert(
|
||||||
|
"supported_content_types".to_string(),
|
||||||
|
serde_json::json!(["text"]),
|
||||||
|
);
|
||||||
|
extra.insert("top_p".to_string(), serde_json::json!(0.9));
|
||||||
|
let provider = make_provider(extra);
|
||||||
|
let filtered = provider.request_model_extra();
|
||||||
|
// 内部字段被过滤
|
||||||
|
assert!(!filtered.contains_key("supported_content_types"));
|
||||||
|
// 业务字段保留
|
||||||
|
assert_eq!(filtered.get("top_p").and_then(|v| v.as_f64()), Some(0.9));
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn test_request_model_extra_empty_when_only_internal() {
|
||||||
|
let mut extra = HashMap::new();
|
||||||
|
extra.insert(
|
||||||
|
"supported_content_types".to_string(),
|
||||||
|
serde_json::json!(["text", "image"]),
|
||||||
|
);
|
||||||
|
let provider = make_provider(extra);
|
||||||
|
assert!(provider.request_model_extra().is_empty());
|
||||||
|
}
|
||||||
|
|
||||||
|
// ---- supports_images: 配置驱动 ----
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn test_supports_images_default_true() {
|
||||||
|
let provider = make_provider(HashMap::new());
|
||||||
|
assert!(provider.supports_images());
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn test_supports_images_disabled_via_config() {
|
||||||
|
let mut extra = HashMap::new();
|
||||||
|
extra.insert(
|
||||||
|
"supported_content_types".to_string(),
|
||||||
|
serde_json::json!(["text"]),
|
||||||
|
);
|
||||||
|
let provider = make_provider(extra);
|
||||||
|
assert!(!provider.supports_images());
|
||||||
|
}
|
||||||
|
|
||||||
|
// ---- AnthropicResponse 反序列化 ----
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn test_deserialize_response_with_text_and_tool_use() {
|
||||||
|
let json = r#"{
|
||||||
|
"id": "msg_001",
|
||||||
|
"model": "claude-3-sonnet",
|
||||||
|
"content": [
|
||||||
|
{"type": "text", "text": "I'll use a tool"},
|
||||||
|
{"type": "tool_use", "id": "call_1", "name": "bash", "input": {"cmd": "ls"}}
|
||||||
|
],
|
||||||
|
"usage": {"input_tokens": 10, "output_tokens": 20}
|
||||||
|
}"#;
|
||||||
|
let resp: AnthropicResponse = serde_json::from_str(json).unwrap();
|
||||||
|
assert_eq!(resp.id, "msg_001");
|
||||||
|
assert_eq!(resp.content.len(), 2);
|
||||||
|
match &resp.content[0] {
|
||||||
|
AnthropicContent::Text { text } => assert_eq!(text, "I'll use a tool"),
|
||||||
|
_ => panic!("expected Text"),
|
||||||
|
}
|
||||||
|
match &resp.content[1] {
|
||||||
|
AnthropicContent::ToolUse { id, name, input } => {
|
||||||
|
assert_eq!(id, "call_1");
|
||||||
|
assert_eq!(name, "bash");
|
||||||
|
assert_eq!(input["cmd"], "ls");
|
||||||
|
}
|
||||||
|
_ => panic!("expected ToolUse"),
|
||||||
|
}
|
||||||
|
assert_eq!(resp.usage.input_tokens, 10);
|
||||||
|
assert_eq!(resp.usage.output_tokens, 20);
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn test_deserialize_response_thinking_variant() {
|
||||||
|
let json = r#"{
|
||||||
|
"id": "msg_002",
|
||||||
|
"model": "claude-3",
|
||||||
|
"content": [
|
||||||
|
{"type": "thinking", "thinking": "internal reasoning"}
|
||||||
|
],
|
||||||
|
"usage": {"input_tokens": 5, "output_tokens": 5}
|
||||||
|
}"#;
|
||||||
|
let resp: AnthropicResponse = serde_json::from_str(json).unwrap();
|
||||||
|
assert_eq!(resp.content.len(), 1);
|
||||||
|
match &resp.content[0] {
|
||||||
|
AnthropicContent::Thinking { thinking } => {
|
||||||
|
assert_eq!(thinking, "internal reasoning");
|
||||||
|
}
|
||||||
|
_ => panic!("expected Thinking"),
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn test_deserialize_response_empty_content() {
|
||||||
|
let json = r#"{
|
||||||
|
"id": "msg_003",
|
||||||
|
"model": "claude-3",
|
||||||
|
"content": [],
|
||||||
|
"usage": {"input_tokens": 1, "output_tokens": 1}
|
||||||
|
}"#;
|
||||||
|
let resp: AnthropicResponse = serde_json::from_str(json).unwrap();
|
||||||
|
assert!(resp.content.is_empty());
|
||||||
|
}
|
||||||
|
|
||||||
|
// ---- format_error_chain ----
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn test_format_error_chain_single() {
|
||||||
|
let err = std::io::Error::new(std::io::ErrorKind::Other, "single error");
|
||||||
|
let chain = format_error_chain(&err);
|
||||||
|
assert_eq!(chain, "single error");
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn test_format_error_chain_nested() {
|
||||||
|
let inner = std::io::Error::new(std::io::ErrorKind::Other, "root cause");
|
||||||
|
let outer = std::io::Error::new(std::io::ErrorKind::Other, "outer error");
|
||||||
|
// std::io::Error 的 source 链需要手动构造,这里用 anyhow 风格的 context 模拟
|
||||||
|
// 简化:验证无 source 时不拼接 "caused by"
|
||||||
|
let chain = format_error_chain(&outer);
|
||||||
|
assert!(chain.contains("outer error"));
|
||||||
|
assert!(!chain.contains("caused by"));
|
||||||
|
let _ = inner; // 抑制未使用警告
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|||||||
Loading…
x
Reference in New Issue
Block a user