从零构建 AI Agent:Tokimo Agent 深度源码剖析
本文以 Tokimo 项目的真实 Rust 实现为例,从 HTTP 协议层到 UI 渲染层,逐层拆解一个 production-grade AI Agent 是如何从零构建的。
目录
- 全景架构
- 与 LLM 对话:HTTP 请求的真相
- 流式响应:SSE 解析的字节级细节
- ReAct 循环:Agent 的心跳
- 工具系统:让 LLM 拥有双手
- 深度思考:Extended Thinking 的三种实现
- 记忆系统:会话记忆与自动压缩
- MCP 协议:可插拔的工具生态
- 事件协议:从 Rust 到浏览器的全链路
- 容错设计:重试、熔断与优雅降级
- 总结:从零开始你需要什么
1. 全景架构
Tokimo Agent 完全用 Rust 实现,分为四个子 crate:
tokimo-agent/
├── packages/core/ # 类型定义、AiProvider trait、三大提供商实现
├── packages/tools/ # Tool trait、结果存储、Bash 只读分析
├── packages/compaction/ # 自动压缩、会话记忆提取
└── packages/runtime/ # AgentRunner、ReAct 循环、工具调度、流式执行
核心数据流:
用户输入
│
▼
┌─────────────┐ ┌──────────────┐ ┌─────────────┐
│ ChatInput │───▶│ SSE Stream │───▶│ AgentRunner │
│ (React) │ │ (HTTP/SSE) │ │ (Rust/tokio)│
└─────────────┘ └──────────────┘ └──────┬──────┘
│
┌─────────▼─────────┐
│ ReAct Loop │
│ ┌───────────────┐ │
│ │ LLM Provider │ │
│ │ (stream_chat) │ │
│ └───────┬───────┘ │
│ │ │
│ ┌───────▼───────┐ │
│ │ Tool Executor │ │
│ │ (concurrent) │ │
│ └───────┬───────┘ │
│ │ │
│ ┌───────▼───────┐ │
│ │ Event Emitter │ │
│ └───────────────┘ │
└────────────────────┘
关键设计决策:Agent 运行时完全在 Rust tokio 异步运行时中执行,通过 SSE(Server-Sent Events)向浏览器推送实时事件。这意味着 LLM 的每一个 token 生成、每一个工具调用,都能以毫秒级延迟到达前端。
2. 与 LLM 对话:HTTP 请求的真相
Agent 的核心是与 LLM 提供商的 HTTP 通信。Tokimo 支持三大提供商,每个都有完全不同的请求格式。
2.1 Anthropic Claude — 最复杂的协议
源码位置:packages/core/src/providers/anthropic.rs
实际发出的 HTTP 请求:
POST https://api.anthropic.com/v1/messages
Content-Type: application/json
x-api-key: sk-ant-xxx
anthropic-version: 2023-06-01
{
"model": "claude-sonnet-4-20250514",
"max_tokens": 8192,
"stream": true,
"system": "You are a helpful assistant...", // ← 系统提示词是顶层字段,不是 message
"messages": [
{
"role": "user",
"content": "帮我写一个排序算法"
},
{
"role": "assistant",
"content": [ // ← Anthropic 用 content blocks 数组
{
"type": "thinking", // ← 思考块(必须带 signature)
"thinking": "用户想要一个排序算法...",
"signature": "ErUBCkYIAxgCIkA..." // ← 原样回传
},
{
"type": "text",
"text": "好的,我来写一个快速排序..."
},
{
"type": "tool_use", // ← 工具调用块
"id": "toolu_01ABC123",
"name": "Bash",
"input": {"command": "python3 sort.py"}
}
]
},
{
"role": "user", // ← 工具结果以 user 角色发送
"content": [
{
"type": "tool_result",
"tool_use_id": "toolu_01ABC123",
"content": "排序完成,输出: [1,2,3,5,8]"
}
]
}
],
"thinking": { // ← 深度思考开关
"type": "enabled",
"budget_tokens": 10000
},
"effort": "high", // ← Opus 4.7 推理强度
"tools": [ // ← 工具定义(从 OpenAI 格式转换)
{
"name": "Bash",
"description": "Execute a shell command",
"input_schema": { // ← Anthropic 用 input_schema,不是 parameters
"type": "object",
"properties": {
"command": {"type": "string"}
},
"required": ["command"]
}
},
{
"type": "web_search_20250305", // ← 内置搜索工具
"name": "web_search"
}
]
}
关键细节:
-
系统提示词的位置:Anthropic 要求
system作为顶层字段,而不是 messages 数组中的一个 message。代码中通过filter(|m| m.role != "system")过滤掉系统消息,然后单独提取:if let Some(sys) = request.messages.iter().find(|m| m.role == "system") { body["system"] = serde_json::json!(sys.content); } -
Content Blocks 协议:Anthropic 不用简单的字符串 content,而是用 content blocks 数组,每个 block 有
type字段(text、thinking、tool_use、tool_result、image)。这意味着一条 assistant 消息可以同时包含思考过程、文本回复和多个工具调用。 -
Thinking Signature:Anthropic 的深度思考块必须带
signature字段回传。如果历史消息中有 thinking 块但丢失了 signature,必须整个 omit 掉,否则 API 会报错:let has_thinking = m.reasoning_content.is_some() && m.thinking_signature.is_some(); // 只有两者都存在时才序列化 thinking block -
工具结果的角色:工具执行结果以
role: "user"发送,content 是tool_result类型的 block。这与其他提供商不同。 - 交替角色约束:Anthropic 要求消息必须 user/assistant 交替。如果第一条消息是 assistant,会自动插入
{"role": "user", "content": "Continue."}。
2.2 OpenAI 兼容协议 — 最广泛的生态
源码位置:packages/core/src/providers/openai_compat.rs
实际发出的 HTTP 请求:
POST https://api.openai.com/v1/chat/completions
Content-Type: application/json
Authorization: Bearer sk-xxx
{
"model": "gpt-4o",
"stream": true,
"stream_options": {"include_usage": true}, // ← 必须显式请求 token 统计
"messages": [
{"role": "system", "content": "You are..."}, // ← 系统消息是普通 message
{"role": "user", "content": "帮我写排序"},
{
"role": "assistant",
"content": null, // ← 有 tool_calls 时 content 为 null
"reasoning_content": "用户想要...", // ← DeepSeek V4 扩展字段
"tool_calls": [
{
"id": "call_abc123",
"type": "function",
"function": {
"name": "Bash",
"arguments": "{\"command\":\"python3 sort.py\"}" // ← arguments 是 JSON 字符串
}
}
]
},
{
"role": "tool", // ← 工具结果用 tool 角色
"tool_call_id": "call_abc123",
"content": "排序完成"
}
],
"reasoning_effort": "high", // ← OpenAI/DeepSeek/Grok 通用
"thinking": {"type": "enabled"}, // ← DeepSeek V4 风格
"verbosity": "medium", // ← GPT-5 系列
"enable_search": true, // ← Qwen 系列搜索开关
"tools": [
{
"type": "function",
"function": {
"name": "Bash",
"description": "Execute a shell command",
"parameters": {
"type": "object",
"properties": {
"command": {"type": "string"}
},
"required": ["command"]
}
}
}
]
}
与 Anthropic 的关键差异:
-
工具调用在 delta 中流式到达:OpenAI 的工具调用参数是分多个 chunk 流式传输的,需要按
index累积。Tokimo 维护了一个tool_call_accumulators向量:// (id, name, arguments, opened, ready) — 5 个状态字段 let mut tool_call_accumulators: Vec<(String, String, String, bool, bool)> = Vec::new(); -
工具名去重问题:OpenAI 规范只在第一个 delta 中发送
function.name,但 xAI、DeepSeek 等兼容提供商会在每个 chunk 中重发完整名字。简单push_str会把"Bash"变成"BashBash"。Tokimo 用merge_tool_call_name函数处理:fn merge_tool_call_name(acc: &mut String, incoming: &str) { if acc == incoming { return; } // 去重:完全相同则跳过 acc.push_str(incoming); // 不同则追加(处理逐字符流式) } -
DeepSeek V4 的 reasoning_content 约束:如果开启 thinking mode,历史中每条 assistant 消息都必须带
reasoning_content字段(即使是空字符串),否则 API 报错。Tokimo 用force_assistant_reasoning标志自动填充:let force_assistant_reasoning = request.reasoning_params.thinking_type.as_deref() == Some("enabled") || request.messages.iter().any(|m| m.role == "assistant" && m.reasoning_content.is_some()); - thinking.type=disabled 与 reasoning_effort 互斥:DeepSeek 不允许同时发送两者,代码中有专门处理。
2.3 Google Gemini — 最独特的格式
源码位置:packages/core/src/providers/google.rs
实际发出的 HTTP 请求:
POST https://generativelanguage.googleapis.com/v1beta/models/gemini-2.5-pro:streamGenerateContent?alt=sse&key=AIzaSyxxx
Content-Type: application/json
{
"systemInstruction": { // ← 系统提示词的特殊格式
"parts": [{"text": "You are..."}]
},
"contents": [ // ← 不叫 messages,叫 contents
{
"role": "user", // ← user 角色相同
"parts": [{"text": "帮我写排序"}]
},
{
"role": "model", // ← 不叫 assistant,叫 model
"parts": [
{
"functionCall": { // ← 不叫 tool_use,叫 functionCall
"name": "Bash",
"args": {"command": "python3 sort.py"} // ← args 是对象,不是 JSON 字符串
},
"thoughtSignature": "xxx" // ← Gemini 2.5+ 的签名
}
]
},
{
"role": "user", // ← 工具结果也是 user 角色
"parts": [
{
"functionResponse": { // ← 不叫 tool_result,叫 functionResponse
"name": "Bash",
"response": {"output": "排序完成"} // ← response 是对象
}
}
]
}
],
"generationConfig": {
"temperature": 0.7,
"maxOutputTokens": 8192,
"thinkingConfig": { // ← 思考配置
"includeThoughts": true,
"thinkingLevel": "high", // ← Gemini 特有的思考级别
"thinkingBudget": 24576 // ← 自动限制为 max_output/2
}
},
"tools": [
{"google_search": {}}, // ← 内置搜索
{
"functionDeclarations": [ // ← 工具定义格式完全不同
{
"name": "Bash",
"description": "Execute a shell command",
"parameters": {
"type": "object",
"properties": {
"command": {"type": "string"}
}
}
}
]
}
],
"toolConfig": {
"includeServerSideToolInvocations": true // ← 搜索+工具共存时必须
}
}
Gemini 的特殊处理:
-
JSON Schema 兼容性:Gemini 不支持 JSON Schema 的某些特性(如
additionalProperties、$ref),需要convert_json_schema_for_google()函数做转换。 -
思考预算自动封顶:Flash 模型上限 24576 tokens,Pro 模型上限 65536 tokens:
let model_thinking_cap: u32 = if request.model.contains("flash") { 24_576 } else { 65_536 }; budget = budget.min(model_thinking_cap); - Grounding Metadata:Gemini 的搜索结果通过
groundingMetadata字段返回,包含搜索查询、引用和支持段落。
3. 流式响应:SSE 解析的字节级细节
所有三个提供商都使用 SSE(Server-Sent Events)协议返回流式数据。Tokimo 的解析器有一个精妙的字节级设计。
3.1 为什么用字节缓冲区而不是字符串?
源码位置:packages/core/src/providers/anthropic.rs:228-268
let mut buffer: Vec<u8> = Vec::new(); // ← 字节缓冲区,不是 String
while let Some(chunk_result) = pinned.next().await {
let chunk = match chunk_result { ... };
buffer.extend_from_slice(&chunk); // ← 追加原始字节
while let Some(pos) = buffer.iter().position(|&b| b == b'\n') {
let drained: Vec<u8> = buffer.drain(..=pos).collect();
let line_bytes = &drained[..drained.len() - 1]; // 去掉 \n
let line = match std::str::from_utf8(line_bytes) {
Ok(s) => s.trim().to_string(),
Err(_) => continue, // ← 不完整的 UTF-8 序列,等下一个 chunk
};
// ... 解析 SSE 行
}
}
为什么这样设计? 源码注释解释得很清楚:
\n是 UTF-8 中的单字节码点(0x0A),永远不可能出现在多字节序列内部。所以每个\n之前的字节总是完整的、有效的 UTF-8 前缀,即使前一个 chunk 在多字节字符中间断开。
这避免了两种常见错误:
- O(n²) 重分配:旧的
buffer = buffer[pos+1..].to_string()模式每次都要复制 - 有损替换:
String::from_utf8_lossy会在字符中间产生U+FFFD替换字符
3.2 Anthropic SSE 事件流
Anthropic 的 SSE 格式是双行结构:
event: message_start
data: {"type":"message_start","message":{"id":"msg_xxx","usage":{"input_tokens":100}}}
event: content_block_start
data: {"type":"content_block_start","index":0,"content_block":{"type":"thinking"}}
event: content_block_delta
data: {"type":"content_block_delta","index":0,"delta":{"type":"thinking_delta","thinking":"让我"}}
event: content_block_delta
data: {"type":"content_block_delta","index":0,"delta":{"type":"thinking_delta","thinking":"思考..."}}
event: content_block_delta
data: {"type":"content_block_delta","index":0,"delta":{"type":"signature_delta","signature":"ErUBCkY..."}}
event: content_block_start
data: {"type":"content_block_start","index":1,"content_block":{"type":"text"}}
event: content_block_delta
data: {"type":"content_block_delta","index":1,"delta":{"type":"text_delta","text":"好的,"}}
event: content_block_delta
data: {"type":"content_block_delta","index":1,"delta":{"type":"text_delta","text":"我来写..."}}
event: content_block_start
data: {"type":"content_block_start","index":2,"content_block":{"type":"tool_use","id":"toolu_xxx","name":"Bash"}}
event: content_block_delta
data: {"type":"content_block_delta","index":2,"delta":{"type":"input_json_delta","partial_json":"{\"command\":"}}
event: content_block_delta
data: {"type":"content_block_delta","index":2,"delta":{"type":"input_json_delta","partial_json":"\"ls\"}"}}
event: content_block_stop
data: {"type":"content_block_stop","index":2}
event: message_delta
data: {"type":"message_delta","delta":{"stop_reason":"tool_use"},"usage":{"output_tokens":150}}
Tokimo 的解析器忽略 event: 行,只处理 data: 行,通过 json["type"] 字段分派:
match event_type {
"message_start" => { /* 捕获 input_tokens + cache_tokens */ }
"content_block_start" if type == "tool_use" => { /* 开始累积工具调用 */ }
"content_block_delta" => {
match delta_type {
"text_delta" => yield ContentDelta,
"thinking" | "thinking_delta" => yield ReasoningDelta,
"signature_delta" => yield ThinkingSignature,
"input_json_delta" => { /* 累积工具参数 */ }
"citations_delta" => { /* 收集引用 */ }
}
}
"content_block_stop" if tool_use => { /* 工具调用完成,立即发射 ToolCallReady */ }
"message_delta" => { /* 捕获 output_tokens,发射 MessageEnd */ }
}
关键设计:ToolCallReady 在 content_block_stop 时立即发射,而不是等到整个 message 结束。这允许流式工具执行器在 API 还在生成后续工具块时就开始执行已完成的工具。
3.3 OpenAI SSE 事件流
OpenAI 的格式更简单,单行 data::
data: {"id":"chatcmpl-xxx","choices":[{"delta":{"role":"assistant"}}]}
data: {"choices":[{"delta":{"reasoning_content":"让我思考"}}]}
data: {"choices":[{"delta":{"reasoning_content":"一下..."}}]}
data: {"choices":[{"delta":{"content":"好的,"}}]}
data: {"choices":[{"delta":{"content":"我来写..."}}]}
data: {"choices":[{"delta":{"tool_calls":[{"index":0,"id":"call_xxx","function":{"name":"Bash"}}]}}]}
data: {"choices":[{"delta":{"tool_calls":[{"index":0,"function":{"arguments":"{\"comm"}}]}}]}
data: {"choices":[{"delta":{"tool_calls":[{"index":0,"function":{"arguments":"and\":\"ls\"}"}}]}}]}
data: {"choices":[{"finish_reason":"tool_calls"}]}
data: {"usage":{"prompt_tokens":100,"completion_tokens":150,"prompt_tokens_details":{"cached_tokens":50}}}
data: [DONE]
工具调用的流式累积是 OpenAI 协议中最复杂的部分。Tokimo 用协议结构(而非 JSON 有效性)判断工具调用是否完成:
// 完备性由协议结构决定,不是尝试解析部分累加器
// (a) 新的更高 index 出现 → 之前的 index 已完成
// (b) finish_reason == "tool_calls" → 所有 open 的 index 已完成
3.4 统一的 ChatStreamEvent
三个提供商的原始 SSE 事件最终被归一化为统一的 ChatStreamEvent 枚举:
pub enum ChatStreamEvent {
MessageStart { message_id, model, provider, created_at },
ContentDelta { delta: String }, // 文本片段
ReasoningDelta { delta: String }, // 思考片段
ToolCallStart { tool_call_id, tool_name, arguments },
ToolCallDelta { tool_call_id, tool_name, partial_json },
ToolCallReady { call: ToolCallRequest }, // 单个工具调用完成
ToolCallsRequested { calls: Vec<...> }, // 批量工具调用(Google 风格)
ToolCallResult { tool_call_id, result },
GroundingMetadata { metadata }, // 搜索引用
ThinkingSignature { signature }, // Anthropic thinking 签名
MessageEnd { usage, reasoning_duration }, // 消息结束 + token 统计
Error,
}
这个归一化层是多提供商支持的关键——上层的 ReAct 循环和工具执行器不需要知道底层是哪个提供商。
4. ReAct 循环:Agent 的心跳
源码位置:packages/runtime/src/runtime/internal/react_loop.rs
ReAct(Reasoning + Acting)循环是 Agent 的核心状态机。Tokimo 的实现是一个单循环,每次迭代做一件事:
┌─────────────────────────────────────────────────────┐
│ ReAct Loop │
│ │
│ 1. 构建请求(首次 or follow-up) │
│ │ │
│ ▼ │
│ 2. 调用 provider.stream_chat() │
│ │ │
│ ▼ │
│ 3. drain_provider() — 流式消费所有事件 │
│ │ │
│ ▼ │
│ 4. finalize_llm_call() — 助手消息入历史 │
│ │ │
│ ├─ 有工具调用? │
│ │ ├─ 是 → collect_tool_results() → continue │
│ │ └─ 否 → 有待注入消息? │
│ │ ├─ 是 → 注入 → continue │
│ │ └─ 否 → break(自然结束) │
│ │ │
│ ▼ │
│ 5. 自动压缩检查 + 会话记忆提取 │
└─────────────────────────────────────────────────────┘
4.1 核心循环代码
loop {
// 1. 构建请求
let stream_result = if is_first_call {
turn.provider.stream_chat(&turn.http, &turn.base_request).await
} else {
let req = turn.build_follow_up_request();
turn.provider.stream_chat(&turn.http, &req).await
};
// 处理 "prompt too long" 错误 → 反应式压缩
let provider_stream = match stream_result {
Ok(s) => s,
Err(e) if Turn::should_attempt_reactive(&e.to_string()) => {
// 强制压缩 → 重试
turn.force_reactive_compact(seed).await
}
Err(e) => { yield Err(e); return; }
};
// 2. 流式消费
let drain = turn.drain_provider(provider_stream);
while let Some(ev) = drain.next().await {
yield Ok(ev); // 转发给前端
}
// 3. 取回检测到的工具调用
let tool_calls = turn.take_detected_tool_calls();
// 4. 无条件将助手消息加入历史
turn.finalize_llm_call(&tool_calls);
// 5. 分支
if !tool_calls.is_empty() {
// 执行工具 → 结果加入历史 → continue
for ev in turn.collect_tool_results(&tool_calls).await {
yield Ok(ev);
}
turn.maybe_auto_compact(estimated).await;
continue;
}
// 无工具调用 → 检查待注入消息
let pending = turn.drain_pending_messages();
if pending.is_empty() {
break; // 真正结束
}
// 注入消息 → continue(同一 run 内继续,用户看到连贯回复)
}
4.2 关键设计:无条件 finalize
这是防止 "essay printed twice" 回归的关键。无论下一步是执行工具、注入待处理消息还是真正结束,助手消息都必须在分支之前加入对话历史。旧架构有三个近乎相同的状态机副本,在 pending 消息到达时会丢失模型的最后回复。
4.3 熔断器
防止工具调用死循环:
const MAX_CONSECUTIVE_ERROR_TURNS: usize = 10;
let mut consecutive_error_turns: usize = 0;
// 每轮结束后检查
if turn.last_turn_all_errors {
consecutive_error_turns += 1;
if consecutive_error_turns >= MAX_CONSECUTIVE_ERROR_TURNS {
yield Err("tool_loop_detected: stopped after N consecutive turns where every tool call failed");
return;
}
} else {
consecutive_error_turns = 0;
}
5. 工具系统:让 LLM 拥有双手
5.1 Tool Trait
源码位置:packages/runtime/src/runtime/tool.rs
#[async_trait]
pub trait Tool: Send + Sync {
fn spec(&self) -> &ToolSpec; // 工具元数据
fn prepare_arguments(&self, raw: serde_json::Value) -> serde_json::Value { raw } // 参数预处理
async fn invoke(&self, ctx: &mut InvokeCtx<'_>) -> ToolOutput; // 执行
}
ToolSpec 包含:
name/description— 告诉 LLM 这个工具是什么parameters— JSON Schema,定义参数格式category— Core / Agent / Plan / Task / Team / Mcp / Experimentalconcurrency— Always(可并行)/ Never(必须串行)/ InputDependent(看参数决定)
5.2 工具定义如何变成 LLM 能理解的格式
所有工具定义存储为 OpenAI function-calling JSON 格式:
{
"type": "function",
"function": {
"name": "Bash",
"description": "Execute a shell command",
"parameters": {
"type": "object",
"properties": {
"command": {"type": "string", "description": "The command to execute"},
"timeout": {"type": "number", "description": "Timeout in ms"}
},
"required": ["command"]
}
}
}
在发送给不同提供商时,Tokimo 自动转换格式:
- Anthropic:去掉
function包装,parameters→input_schema - Gemini:转换为
functionDeclarations数组,处理 JSON Schema 兼容性 - OpenAI:直接使用
5.3 流式工具执行器
源码位置:packages/runtime/src/runtime/internal/streaming_tool_executor.rs
这是 Tokimo 最精妙的设计之一。工具不是等所有工具调用都到齐后才执行,而是边流式接收边执行:
struct TrackedTool {
tc: ToolCallRequest,
status: ToolStatus, // Queued → Executing → Completed → Yielded
is_concurrency_safe: bool,
result_slot: Arc<Mutex<Option<ToolExecResult>>>,
handle: Option<tokio::task::JoinHandle<()>>, // 可以 abort
}
并发控制:
fn can_execute(&self, is_concurrency_safe: bool) -> bool {
let executing = self.tools.iter().filter(|t| t.status == ToolStatus::Executing);
executing.is_empty()
|| (is_concurrency_safe && executing.iter().all(|t| t.is_concurrency_safe))
}
规则:
- 没有正在执行的工具 → 可以执行任何工具
- 有正在执行的工具 → 只有标记为
concurrency_safe的工具可以并行执行 - 写操作(
ConcurrencyMode::Never)必须独占
生命周期管理:当 executor 被 drop 或 reset 时,所有正在执行的工具任务会被 abort():
impl Drop for StreamingToolExecutor {
fn drop(&mut self) {
for tool in &mut self.tools {
if let Some(h) = tool.handle.take() {
h.abort(); // 取消正在执行的工具
}
}
}
}
这确保了用户点击"停止生成"后,正在运行的 Bash 命令、HTTP 请求等都会被立即取消。
5.4 Bash 只读分析
Tokimo 有一个静态分析器判断 Bash 命令是否只读(可以安全并行执行):
packages/tools/src/bash_readonly/
├── mod.rs # is_bash_read_only() 入口
├── allowlist.rs # 允许的只读命令列表
├── tokenizer.rs # 命令分词器
└── flag_validation.rs # 标志验证
例如 cat file.txt 和 ls -la 是只读的,可以与其它工具并行;rm file.txt 和 echo "x" > file 不是。
5.5 完整的工具列表
Tokimo 注册了 20+ 个内置工具:
| 工具名 | 类别 | 并发 | 功能 |
|---|---|---|---|
| Bash | Core | InputDependent | Shell 命令执行 |
| Read | Core | Always | 读取文件 |
| Write | Core | Never | 写入文件 |
| Edit | Core | Never | 编辑文件(精确替换) |
| Glob | Core | Always | 文件模式匹配搜索 |
| Grep | Core | Always | 内容搜索 |
| LSP | Core | Always | 语言服务器交互 |
| WebFetch | Core | Always | 抓取网页 |
| WebSearch | Core | Always | 搜索引擎 |
| Agent | Agent | Always | 派生子 agent |
| DispatchAgent | Agent | Always | 派生独立 agent |
| Plan | Plan | Always | 进入计划模式 |
| Task | Task | Always | 任务管理 |
| Skill | Agent | Always | 技能执行 |
| Cron | Experimental | Always | 定时任务 |
| Workflow | Agent | Always | 工作流编排 |
| Worktree | Core | Never | Git worktree 管理 |
| AskUserQuestion | Core | Never | 向用户提问 |
| Memory | Core | Always | 记忆读写 |
| ToolSearch | Core | Always | 搜索可用工具 |
| Config | Core | Always | 配置管理 |
6. 深度思考:Extended Thinking 的三种实现
深度思考(Extended Thinking)让 LLM 在回答前进行内部推理。三家提供商的实现完全不同。
6.1 Anthropic 实现
{
"thinking": {
"type": "enabled", // enabled / disabled / adaptive
"budget_tokens": 10000 // 思考 token 预算
}
}
响应中的思考内容是独立的 content block:
{
"type": "content_block_delta",
"delta": {
"type": "thinking_delta", // ← 思考片段
"thinking": "让我分析这个问题..."
}
}
思考完成后会收到 signature_delta,这个签名必须在后续请求中原样回传。
6.2 OpenAI/DeepSeek 实现
{
"reasoning_effort": "high", // minimal / low / medium / high / xhigh
"thinking": {"type": "enabled"} // DeepSeek V4 风格
}
响应中的思考内容在 reasoning_content 字段:
{
"choices": [{
"delta": {
"reasoning_content": "让我分析...", // ← 思考内容
"content": "好的,我来..." // ← 正式回复
}
}]
}
DeepSeek 约束:thinking.type=disabled 和 reasoning_effort 互斥;如果历史中有 reasoning_content,后续所有 assistant 消息都必须带这个字段。
6.3 Google Gemini 实现
{
"generationConfig": {
"thinkingConfig": {
"includeThoughts": true,
"thinkingLevel": "high", // minimal / low / medium / high
"thinkingBudget": 24576 // 自动封顶
}
}
}
响应中的思考内容是 thought: true 的 parts:
{
"candidates": [{
"content": {
"parts": [
{"text": "让我思考...", "thought": true}, // ← 思考
{"text": "好的,我来..."} // ← 正式回复
]
}
}]
}
6.4 统一的 ReasoningParams
pub struct ReasoningParams {
pub effort: Option<String>, // OpenAI/Anthropic 通用
pub thinking_type: Option<String>, // Anthropic/DeepSeek
pub thinking_budget: Option<u32>, // Anthropic/Gemini
pub thinking_level: Option<String>, // Gemini
pub text_verbosity: Option<String>, // GPT-5
pub enabled_search: Option<bool>, // Qwen
pub disable_context_caching: Option<bool>,
pub url_context: Option<bool>, // Gemini
pub image_aspect_ratio: Option<String>,
pub image_resolution: Option<String>,
}
前端根据模型的 extend_params 配置显示不同的控件(滑块、开关、下拉框),所有值最终映射到这个统一结构。各提供商只取自己理解的字段,忽略其他。
7. 记忆系统:会话记忆与自动压缩
7.1 自动压缩(Auto-Compaction)
源码位置:packages/compaction/src/auto_compact.rs
当对话 token 数接近上下文窗口限制时,自动触发压缩:
阈值计算:
fn auto_compact_threshold(context_window: Option<u32>, max_output: Option<u32>) -> u32 {
let cw = context_window.unwrap_or(200_000).min(200_000); // 硬上限 200k
let mo = max_output.unwrap_or(8_000);
let reserved = mo.min(20_000); // 为摘要生成预留
let effective = cw - reserved;
effective - 13_000 // 安全缓冲
}
例如 200k 上下文窗口:200000 - 8000 - 13000 = 179000 tokens 时触发压缩。
Token 估算:使用 chars / 4 启发式,锚定到上次 API 返回的 usage:
fn estimate_token_count(messages, last_usage, messages_since_usage) -> u32 {
if let Some(usage) = last_usage {
let anchor = usage.input_tokens + usage.output_tokens;
let new_estimate = messages_since_usage.iter().map(|m| m.content.len() / 4 + 4).sum();
anchor + new_estimate
} else {
messages.iter().map(|m| m.content.len() / 4 + 4).sum()
}
}
熔断器:连续 3 次压缩失败后停止尝试(防止浪费 API 调用)。
反应式压缩:如果 API 返回 "prompt too long" 错误,ReAct 循环会立即触发压缩并重试:
if Turn::should_attempt_reactive(&err_str) {
match turn.force_reactive_compact(seed).await {
Ok((new_msgs, info)) => {
yield Ok(compaction_event(...));
// 用压缩后的消息重试
}
}
}
7.2 会话记忆(Session Memory)
源码位置:packages/compaction/src/session_memory/
会话记忆是一种增量式的对话摘要,在后台持续更新:
触发条件:
- 距上次提取的 token 数超过阈值
- 距上次提取的工具调用数超过阈值
提取流程:
pub async fn extract_session_memory(
extractor: Arc<dyn SessionMemoryExtractor>,
conversation_msgs: Vec<Message>,
state: Arc<Mutex<SessionMemoryState>>,
current_tokens: u32,
) -> Result<u64, String> {
// 1. 标记提取开始
state.lock().await.mark_extraction_started();
// 2. 读取现有记忆(或使用模板)
let current_notes = match tokio::fs::read_to_string(&memory_path).await {
Ok(content) if !content.trim().is_empty() => content,
_ => SESSION_MEMORY_TEMPLATE.to_string(),
};
// 3. 调用提取器(实际实现是一个子 agent + Edit 工具循环)
extractor.extract(conversation_msgs, current_notes, &memory_path).await
}
提取器实现(在 server crate 中):实际是一个 forked agent,使用 Edit 工具增量更新记忆文件。这和 Claude Code 的 runForkedAgent 是同样的设计。
钩子机制:
// hooks_builtin.rs
pub struct SessionMemoryExtractHook;
impl Hook for SessionMemoryExtractHook {
async fn handle(&self, event: HookEvent) {
if let HookEvent::TurnEnd { ctx, .. } = event {
if ctx.should_extract_session_memory() {
tokio::spawn(extract_session_memory(...)); // 后台提取
}
}
}
}
7.3 知识库(Wiki)
知识库是持久化的文档存储,使用 git 做版本控制:
packages/rust-server/src/apps/agent/modules/wiki/
├── content.rs # 内容管理
├── file_repo.rs # 文件仓库
├── run_repo.rs # 运行仓库
└── sandbox.rs # 沙箱挂载
知识库内容可以被挂载到 agent 的沙箱中,让 agent 在工具调用时访问。
8. MCP 协议:可插拔的工具生态
8.1 MCP 是什么
MCP(Model Context Protocol)是一个标准化的工具服务器协议。Tokimo 作为 MCP 客户端,可以连接任意 MCP 服务器来扩展 agent 的能力。
8.2 JSON-RPC 2.0 通信
源码位置:packages/tokimo-package-mcp/src/client.rs
MCP 使用 JSON-RPC 2.0 协议:
// 请求
{"jsonrpc": "2.0", "id": 1, "method": "tools/list", "params": {}}
// 响应
{"jsonrpc": "2.0", "id": 1, "result": {"tools": [
{
"name": "read_file",
"description": "Read a file from disk",
"inputSchema": {
"type": "object",
"properties": {"path": {"type": "string"}},
"required": ["path"]
},
"annotations": {
"readOnlyHint": true,
"destructiveHint": false
}
}
]}}
// 通知(服务器主动推送)
{"jsonrpc": "2.0", "method": "notifications/tools/list_changed"}
8.3 传输层
pub trait McpTransport: Send + Sync {
async fn send(&self, msg: serde_json::Value) -> Result<(), String>;
async fn recv(&self) -> Option<serde_json::Value>;
async fn close(&self);
}
两种实现:
- stdio:启动子进程,通过 stdin/stdout 通信
- http:HTTP POST + SSE 流
8.4 连接管理
pub struct McpConnectionManager {
connections: DashMap<Uuid, Arc<McpConnection>>, // 并发安全的连接池
}
impl McpConnectionManager {
// 后台预热:启动时并发连接所有启用的 MCP 服务器
async fn warmup(&self, configs: &[McpServerConfig]) { ... }
// 将 MCP 工具转换为 agent 的 Tool trait 对象
async fn build_tools_for_agent(&self) -> Vec<Arc<dyn Tool>> { ... }
}
8.5 MCP 工具适配器
MCP 工具通过 McpToolAdapter 桥接到 agent 的工具系统:
struct McpToolAdapter {
connection: Arc<McpConnection>,
tool: McpTool, // name, description, input_schema, annotations
}
#[async_trait]
impl Tool for McpToolAdapter {
fn spec(&self) -> &ToolSpec { ... }
async fn invoke(&self, ctx: &mut InvokeCtx<'_>) -> ToolOutput {
let result = self.connection.call_tool(&self.tool.name, ctx.arguments()).await;
ToolOutput::ok(result.content)
}
}
这意味着 agent 可以无缝调用 MCP 服务器提供的工具,就像调用内置工具一样。
9. 事件协议:从 Rust 到浏览器的全链路
9.1 Rust 端的 Event 枚举
源码位置:packages/runtime/src/runtime/event.rs
pub enum Event {
Conversation { conversation_id, status }, // 会话生命周期
Turn { ref_id, status }, // 一轮对话的生命周期
Message { id, ref_id, role, msg_type, status, payload, model, provider },
Usage { input_tokens, output_tokens, cache_tokens, model, provider },
Compaction { before_tokens, after_tokens, strategy },
SaveContext(ConversationState), // 内部事件
AwaitingUserInput(AwaitingUserInput), // 内部事件
}
Message 的生命周期:
Start → Streaming(delta) → Streaming(delta) → ... → End(content, metadata)
每个阶段都是一个独立的 Event,通过 SSE 推送到前端。
9.2 从 Event 到 SSE
服务端有一个 Projector 层:
Event Stream → Projector → DB 持久化 + SSE 广播
ChatProjector:将 Event 写入数据库(PostgreSQL)ProjectorSink:实现AgentEventSinktraitStreamingBuffers:管理每个连接的 delta 缓冲区
9.3 前端 SSE 事件处理
源码位置:packages/web/src/apps/agent/lib/sse-event-handler.ts
前端通过 EventSource(或 fetch + ReadableStream)接收 SSE 事件,然后应用到 React Query 缓存:
type AgentEvent =
| ConversationEvent // { kind: "conversation", status: "start" | "end" | ... }
| TurnEvent // { kind: "turn", ref_id, status: "start" | "end" }
| MessageEvent // { kind: "message", id, status: "start" | "streaming" | "end" }
| UsageEvent // { kind: "usage", input_tokens, output_tokens, ... }
| CompactionEvent; // { kind: "compaction", before_tokens, after_tokens, ... }
Delta 缓冲:前端使用 requestAnimationFrame 批量刷新 streaming delta,避免每个 token 都触发 React 重渲染:
// 使用 RAF 调度批量刷新
if (!rafScheduled) {
rafScheduled = true;
requestAnimationFrame(() => {
flushPendingDeltas();
rafScheduled = false;
});
}
增量 JSON 解析:工具参数在流式到达时就能部分解析显示:
function tryParsePartialJson(raw: string): JsonValue | null {
// 尝试解析不完整的 JSON,用于实时显示工具参数
}
9.4 完整的数据流
1. 用户输入 "帮我写排序"
│
2. 前端 POST /api/apps/agent/chat
│
3. 服务端 spawn_agent_run() 组装基础设施
│
4. AgentRunner::run() 启动 tokio task
│
5. ReAct 循环调用 provider.stream_chat()
│
6. Anthropic API 返回 SSE 流
│
7. anthropic.rs 解析 SSE → ChatStreamEvent
│
8. turn.rs 归一化 → Event
│
9. Projector 写 DB + 推 SSE
│
10. 前端 EventSource 接收
│
11. sse-event-handler.ts 更新 React Query 缓存
│
12. React 组件重渲染
│
13. 用户看到流式输出
10. 容错设计:重试、熔断与优雅降级
10.1 重试层
源码位置:packages/core/src/providers/retry.rs
RetryingAiProvider 装饰器包装任何 AiProvider:
pub struct RetryingAiProvider {
inner: Box<dyn AiProvider>,
config: RetryConfig,
}
重试策略:
| 条件 | 行为 |
|---|---|
| HTTP 408/429/5xx | 重试 |
| 连接错误 | 重试 |
| 流打开超时(15s) | 重试 |
| HTTP 4xx(非 408/429) | 不重试(永久错误) |
| 流内部错误 | 不重试(会破坏消息历史) |
退避阶梯:3s → 10s → 30s → 1min → 2min,±25% 抖动。
为什么不在流内部重试? 源码注释解释:
By the time the stream is yielding, downstream code has already emitted a
MessageStartevent and may have persisted tokens — replaying the request would double-speak and corrupt message history.
10.2 熔断器
- 工具循环熔断:连续 10 轮全部工具失败 → 强制停止
- 自动压缩熔断:连续 3 次压缩失败 → 停止尝试
- 上下文窗口硬上限:即使模型声明 1M 上下文,也按 200k 处理
10.3 优雅降级
- 模型列表获取失败 → 使用硬编码的默认模型列表
- 会话记忆提取失败 → 释放标志位,下次重试
- MCP 服务器连接失败 → 跳过该服务器的工具,不阻塞 agent
- 工具参数验证失败 → 返回
InputValidationError给模型,让它重试
11. 总结:从零开始你需要什么
如果你想从零构建一个类似的 AI Agent,以下是最低限度的技术栈:
11.1 核心组件
| 组件 | 你需要实现 | 复杂度 |
|---|---|---|
| LLM Provider | HTTP 客户端 + SSE 解析 | ★★★ |
| 消息格式转换 | 每个提供商一套序列化逻辑 | ★★ |
| ReAct 循环 | 状态机:调用 → 检测工具 → 执行 → 循环 | ★★★ |
| 工具系统 | Tool trait + JSON Schema 验证 + 并发控制 | ★★★ |
| 流式输出 | SSE 推送 + 前端 delta 缓冲 | ★★ |
| 错误处理 | 重试 + 熔断 + 降级 | ★★ |
11.2 最小可行实现
如果只支持一个提供商(比如 OpenAI),最小实现大约需要:
- HTTP 客户端(~50 行):发 POST 请求,带 Bearer token
- SSE 解析器(~50 行):解析
data:行,提取 JSON - 消息格式化(~30 行):将对话历史转为 OpenAI messages 数组
- 工具定义(~20 行/工具):JSON Schema + 执行函数
- ReAct 循环(~100 行):核心状态机
- 流式输出(~50 行):SSE 推送到前端
总计约 300 行核心代码就能实现一个基本的 agent。
11.3 从 MVP 到 Production 的差距
Tokimo 的实现大约是 5000+ 行 Rust 代码(不含工具实现),多出来的主要是:
- 多提供商支持:三套完全不同的协议适配
- 流式工具执行:边接收边执行,并发控制
- 上下文管理:自动压缩、会话记忆、token 估算
- 容错:重试、熔断、优雅降级
- MCP 集成:可插拔工具生态
- 前端渲染:20+ 种工具的专用渲染器
11.4 关键架构决策
- Rust + tokio:高并发工具执行需要真正的异步运行时,JS 的 event loop 在 CPU 密集型工具(如文件搜索)上会阻塞
- SSE 而非 WebSocket:SSE 更简单、HTTP 原生支持、自动重连
- 统一的 ChatStreamEvent:归一化层让上层代码不关心提供商差异
- 流式工具执行:不等所有工具到齐,已完成的立即执行
- Arc
:消息用 Arc 包装,per-turn clone 是 O(n) 指针复制 - 字节级 SSE 解析:避免 UTF-8 边界问题和 O(n²) 重分配
本文基于 Tokimo 项目的真实源码分析。所有代码片段都来自实际的生产代码。
版权属于:一名宅。
本文链接:https://zhaiyiming.com/archives/76.html
转载时须注明出处及本声明