From 832ebf266520da38d8d6cca8721aa16a96e5a502 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=E5=BE=90=E6=B6=9B?= Date: Tue, 30 Jun 2026 22:22:12 +0800 Subject: [PATCH] =?UTF-8?q?docs(phase0):=20=E6=96=B0=E5=A2=9E=20Phase=200?= =?UTF-8?q?=20=E5=AE=9E=E6=96=BD=E8=AE=A1=E5=88=92=E6=96=87=E6=A1=A3?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 定义类型层落地计划,涵盖新类型定义、Trait 签名切换及上游适配方案 --- docs/10a-phase0-types-and-trait.md | 997 +++++++++++++++++++++ docs/10b-phase1-provider-adaptation.md | 1133 ++++++++++++++++++++++++ docs/10c-phase2-llm-cycle-simplify.md | 1049 ++++++++++++++++++++++ 3 files changed, 3179 insertions(+) create mode 100644 docs/10a-phase0-types-and-trait.md create mode 100644 docs/10b-phase1-provider-adaptation.md create mode 100644 docs/10c-phase2-llm-cycle-simplify.md diff --git a/docs/10a-phase0-types-and-trait.md b/docs/10a-phase0-types-and-trait.md new file mode 100644 index 0000000..686a0a9 --- /dev/null +++ b/docs/10a-phase0-types-and-trait.md @@ -0,0 +1,997 @@ +# Phase 0 实施计划:类型层落地 + Trait 签名切换 + +> **所属方案**:[10-llm-provider-refinement.md](10-llm-provider-refinement.md) +> +> **前置条件**:无(首个实施阶段) +> +> **产出依赖**:Phase 1(Provider 适配)依赖本阶段的类型定义和 trait 签名 + +--- + +## 目标 + +新增新类型系统 + 切换 `LlmProvider` trait 签名,使全链路使用新类型。Phase 0 结束时 `cargo test` 全部通过。 + +## 原则 + +1. **新类型定义放入新文件**(`message.rs`、`request_v2.rs`、`response_v2.rs`),不堆积到已有类型文件 +2. 已有的 `request.rs`(`OpenaiChatRequest`)、`response.rs`(`OpenaiChatResponse`)、`stream.rs`(旧 `StreamEvent`)**保留原样**,后续 Provider 实现可能作为内部转换目标继续引用 +3. `LlmProvider` trait 签名由 `chat(ChatRequest) → ChatResponse` 切换为 `chat(MessageRequest) → MessageResponse`,**在同一个 Phase 内完成** +4. trait 签名变更导致的编译错误(`StubProvider`、`LlmCycle` 调用点)**在 Phase 0 内全部修复**,不留到 Phase 1 +5. `AgentSession` 等上游中对 `LlmCycle.submit()` 返回值的引用同步适配 +6. 不使用任何新依赖,只在现有 crate 范围内完成 + +--- + +## 涉及文件 + +| 操作 | 文件 | 说明 | +|------|------|------| +| 新增 | `src/llm/types/message.rs` | Message 扁平大枚举 + ContentBlock + 辅助类型 | +| 新增 | `src/llm/types/request_v2.rs` | MessageRequest + ExtraError + extra 访问方法 | +| 新增 | `src/llm/types/response_v2.rs` | MessageResponse + StreamEvent(新) + PartialMessageResponse | +| 追加 | `src/llm/types/mod.rs` | 追加 `pub mod` 声明和 `pub use` 重导出 | +| 修改 | `src/llm/provider.rs` | LlmProvider trait 签名切换 | +| 修改 | `src/agent/builder.rs` | StubProvider 适配新 trait 签名 | +| 修改 | `src/llm/cycle.rs` | build_request/submit/submit_stream/submit_messages/submit_request 适配 | +| 修改 | `src/llm/cycle/retry.rs` | 如有对新 LlmError 类型的引用,同步适配 | +| 修改 | `src/agent/session.rs` | ChatResponse → MessageResponse 引用适配 | +| 修改 | `src/memory/conversation.rs` | 如有对旧类型别名的引用,同步适配 | + +## 实施前基线确认 + +实施者在开始 Phase 0 前应确认: +1. `git status` — 工作区干净,无未提交的修改 +2. `cargo build` — 编译通过,无已有错误 +3. `cargo test` — 所有测试通过 +4. `grep -r "pub type Message = OpenaiChatMessage" src/` — 确认旧类型别名的引用面(供命名冲突处理参考) + +如果基线已有问题,在修复基线后再开始 Phase 0,以免干扰对 Phase 0 改动的判断。 + +--- + +## 任务依赖关系 + +``` +任务 1-4(新类型定义) ← 可并行 + │ + ├──→ 任务 5(新类型单元测试)← 可并行,不阻塞下游 + │ + └──→ 任务 6(trait 签名切换) + │ + └──→ 任务 7(StubProvider 适配) + │ + └──→ 任务 8(LlmCycle 适配) + │ + └──→ 任务 9(上游适配) +``` + +**关键路径**:任务 1 → 6 → 7 → 8 → 9 +**可并行**:任务 5 可与任务 6-9 并行编写,但在执行任务 6 前需确认类型定义已就绪 + +--- + +## 任务 1:新增 `src/llm/types/message.rs` + +定义 IR 层消息类型,基于 10 号文档 §2.1 Decision-01。 + +### 类型定义 + +```rust +use crate::llm::types::shared::ImageDetail; + +/// 跨 Provider 统一的消息类型(扁平大枚举)。 +#[derive(Debug, Clone, Serialize, Deserialize)] +#[serde(rename_all = "snake_case")] +pub enum Message { + System { content: Vec }, + User { content: Vec }, + UserImage { data: String, mime_type: String, detail: ImageDetail }, + Assistant { content: Vec }, + ToolResult { tool_call_id: String, content: Vec, is_error: bool }, +} +``` + +**注意**: +- 9b 文档中原有 `Tool` 变体(对应 OpenAI `tool` role),本设计改为 `ToolResult` +- `ImageDetail` 类型已在 `shared.rs` 中定义(`enum { Auto, Low, High }`,已派生 `Serialize/Deserialize`),此处直接引用 + +### ContentBlock 类型 + +从 9b §3.1 移植,参考 10 号文档 §4 任务 2: + +```rust +#[derive(Debug, Clone, Serialize, Deserialize)] +#[serde(rename_all = "snake_case")] +pub enum ContentBlock { + Text { text: String }, + Image { source: ImageSource }, + Audio { source: AudioSource }, + File { source: FileSource }, + ToolUse { id: String, name: String, input: serde_json::Value }, + ToolResult { tool_use_id: String, content: Vec, is_error: bool }, + Thinking { text: String, signature: Option }, + Extension { kind: String, data: serde_json::Value }, +} +``` + +### ContentBlockType(用于 StreamEvent.ContentBlockStart) + +```rust +#[derive(Debug, Clone, Serialize, Deserialize)] +#[serde(rename_all = "snake_case")] +pub enum ContentBlockType { + Text, + Thinking, + Refusal, + ToolUse { id: String, name: String }, +} +``` + +### 辅助类型 + +从 9b §3.1 移植: + +```rust +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct ImageSource { + pub data: String, + pub mime_type: String, + pub is_url: bool, +} + +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct AudioSource { + pub data: String, + pub format: AudioFormat, // 复用现有 shared.rs 中的 AudioFormat(已派生 Serialize/Deserialize) +} + +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct FileSource { + pub data: String, + pub filename: Option, + pub mime_type: Option, +} +``` + +### 便捷构造函数 + +基于 10 号文档 §2.1 设计: + +```rust +impl Message { + pub fn user_text(text: impl Into) -> Self { ... } + pub fn user_image(data: impl Into, mime_type: impl Into, detail: ImageDetail) -> Self { ... } + pub fn assistant(text: impl Into) -> Self { ... } + pub fn system(text: impl Into) -> Self { ... } + pub fn tool_result(tool_call_id: impl Into, text: impl Into, is_error: bool) -> Self { ... } +} +``` + +**注意**:这里的 `tool_result` 签名与 9b 文档不同——增加了 `is_error: bool` 参数,与扁平大枚举的 `ToolResult` 变体一致。 + +### 关于 serde 策略 + +所有新类型统一规则: +- 使用 `#[derive(Serialize, Deserialize)]`(serde 已是项目依赖) +- 枚举使用 `#[serde(rename_all = "snake_case")]`(与现有项目中 `AudioFormat`、`ImageDetail` 等保持一致) +- 结构体使用默认命名(不额外标注 rename) +- `ContentBlock` 等复合枚举使用外部标记格式(externally tagged,serde 默认行为),后续如需调整 tagging 策略(如 internally tagged)可在实施时按需修改 + +### 验证点 + +- 所有枚举变体 match 穷举性验证(编译器保证) +- `Message` 实现 `Debug + Clone + Serialize + Deserialize` +- `ContentBlock` 实现 `Debug + Clone + Serialize + Deserialize` + +--- + +## 任务 2:移植 ContentBlock 辅助类型 + +### ImageSource + +```rust +pub struct ImageSource { + pub data: String, // base64 或 URL + pub mime_type: String, // "image/png", "image/jpeg", "image/webp" + pub is_url: bool, // true = URL, false = base64 +} +``` + +### AudioSource + +复用 `shared.rs` 中的 `AudioFormat`: + +```rust +pub struct AudioSource { + pub data: String, + pub format: AudioFormat, +} +``` + +### FileSource + +```rust +pub struct FileSource { + pub data: String, + pub filename: Option, + pub mime_type: Option, +} +``` + +### 与 10 号文档的差异说明 + +10 号文档 §4 任务 2 要求从 `9b` 移植 `ImageSource`、`AudioSource`、`FileSource`。其中 `ImageSource` 由 `ImageURL` 简化而来(去掉 `detail` 字段,增加 `is_url` 标记)。`AudioSource` 复用现有 `AudioFormat`。`FileSource` 与现有 `FileData` 结构相似但字段名简化。 + +--- + +## 任务 3:新增 `src/llm/types/request_v2.rs` + +基于 9b §3.6 MessageRequest + §3.7 ExtraError。 + +### MessageRequest + +```rust +use std::collections::HashMap; +use crate::llm::types::request::ToolChoice; +// ToolDefinition 已在 types/mod.rs 中定义为 pub type ToolDefinition = OpenaiToolDefinition + +#[derive(Debug, Clone, Default, Serialize, Deserialize)] +pub struct MessageRequest { + pub model: String, + pub messages: Vec, + pub tools: Vec, + pub tool_choice: ToolChoice, + pub max_tokens: Option, + pub temperature: Option, + pub top_p: Option, + pub stop_sequences: Vec, + pub stream: bool, + pub thinking: Option, + pub extra: HashMap, +} +``` + +**说明**: +- `system` 字段不存在——系统提示通过 `Message::System { content }` 在 `messages` 列表中表达 +- `ToolDefinition` 复用现有的 `OpenaiToolDefinition` 类型别名(`src/llm/types/mod.rs` 中已有 `pub type ToolDefinition = OpenaiToolDefinition`;如果 `message.rs` 和 `request_v2.rs` 不在 `types/` module 内,通过 `use crate::llm::types::ToolDefinition;` 引入) +- `ToolChoice` 复用现有的 `ToolChoice` 枚举(`src/llm/types/request.rs`;通过 `use crate::llm::types::request::ToolChoice;` 引入) +- `#[derive(Default)]` 确保 `..Default::default()` 可用(如 `build_request` 中使用) +- `MessageRequest` 要求 `Message` 满足 `Serialize + Deserialize`(任务 1 已派生) + +### ⚠️ 命名冲突:旧 `type Message` 与新 `enum Message` + +`types/mod.rs` 中现有 `pub type Message = OpenaiChatMessage;` 别名。Phase 0 新增 `message.rs` 并导出 `pub enum Message` 后,在 `types::` 命名空间下产生重定义冲突。 + +**解决方案(三选一,推荐选项 A):** + +| 选项 | 操作 | 影响 | +|------|------|------| +| **A(推荐)** | 在 `types/mod.rs` 中移除 `pub type Message = OpenaiChatMessage;`,所有仍引用 `types::Message` 的地方改为直接使用 `OpenaiChatMessage` | 旧别名已不必要(新代码都引用 `Message`),移除后无外部使用者依赖此别名 | +| B | 将旧别名重命名为 `pub type OpenaiMessage = OpenaiChatMessage;` | 需同步更新所有引用点,Phase 2 清理时再删除 | +| C | `message.rs` 中不直接 `pub use` `Message`,而是通过 `MessageEnum` 等中间名称导出 | 对外接口不干净,不推荐 | + +**建议 Phase 0 实施时采用选项 A**,因为 `types::Message` 作为 `OpenaiChatMessage` 的别名在现有代码中引用面很小(`grep "types::Message" src/ -r` 确认),且 Phase 2 最终会完全淘汰 `OpenaiChatMessage`。 + +### ThinkingConfig + +```rust +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct ThinkingConfig { + pub budget_tokens: u32, +} +``` + +### ExtraError 枚举 + +基于 9b §3.7 ExtraError: + +```rust +#[derive(thiserror::Error, Debug)] +pub enum ExtraError { + #[error("extra 字段 `{key}` 类型不匹配: {details}")] + TypeMismatch { key: String, details: String }, + #[error("extra 反序列化失败: {0}")] + Deserialize(String), +} +``` + +### extra 访问方法 + +```rust +impl MessageRequest { + pub fn get_extra(&self, key: &str) -> Result, ExtraError> { ... } + pub fn get_extra_opt(&self, key: &str) -> Option { ... } + pub fn get_extra_as(&self) -> Result { ... } + pub fn set_extra(&mut self, key: impl Into, value: impl Into) { ... } +} +``` + +### 验证点 + +- `MessageRequest` 可构造 +- `set_extra` / `get_extra` 基本路径:设置后能正确读取 +- `get_extra` 类型不匹配时返回 `Err(ExtraError::TypeMismatch)` + +--- + +## 任务 4:新增 `src/llm/types/response_v2.rs` + +基于 10 号文档 §2.3 Decision-03 修订后的定义。 + +### StopReason 枚举 + +基于 9b §3.3: + +```rust +#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)] +#[serde(rename_all = "snake_case")] +pub enum StopReason { + Stop, Length, ToolUse, ContentFilter, MaxTokens, StopSequence, Other, +} +``` + +### MessageResponse + +```rust +use std::collections::HashMap; +use crate::llm::types::Usage; + +/// 复用现有 Usage 类型(已在 types/mod.rs 中定义,已派生 Serialize/Deserialize)。 +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct MessageResponse { + pub id: String, + pub model: String, + pub message: Message, + pub usage: Usage, + pub stop_reason: StopReason, + pub extra: HashMap, +} + +impl MessageResponse { + pub fn text(&self) -> String { ... } +} +``` + +### PartialUsage + +```rust +#[derive(Debug, Clone, Default, Serialize, Deserialize)] +pub struct PartialUsage { + pub prompt_tokens: Option, + pub completion_tokens: Option, + pub total_tokens: Option, + pub completion_tokens_details: Option, + pub prompt_tokens_details: Option, +} +``` + +### StreamEvent(高精度版) + +基于 10 号文档 §2.3 Decision-03: + +```rust +#[derive(Debug, Clone, Serialize, Deserialize)] +#[serde(rename_all = "snake_case")] +pub enum StreamEvent { + // Meta + MessageStart { id: String, model: String }, + // Content Block 边界 + ContentBlockStart { index: u32, block_type: ContentBlockType }, + ContentBlockEnd { index: u32 }, + // 块内增量 + TextDelta { text: String }, + ThinkingDelta { text: String }, + RefusalDelta { text: String }, + ToolCallArgumentsDelta { index: u32, arguments: String }, + ToolCallEnd { index: u32 }, + // 汇总 + CostUpdate { usage: PartialUsage }, + /// 消息完成——唯一可靠的完整响应来源。full_response 携带完整的 MessageResponse。 + MessageComplete { full_response: MessageResponse }, + // 错误 + Error { message: String }, +} +``` + +**说明**:与 9c 原有设计不同,`MessageComplete` 不再携带独立的 `stop_reason` 和 `thinking_signature` 字段——这些信息已在 `full_response` 中。 + +### PartialMessageResponse + +基于 9c §4.4,包含 `ContentBlockBuilder` 定义和 `apply_to` / `finalize` 方法。 + +**关键变更**(对应 10 号文档 §2.3 修订): +- `thinking_signature` 改为由 Provider 直接调用 `set_thinking_signature()` 写入内部状态,不再经过事件层 +- `MessageComplete.apply_to` 不再处理 `stop_reason` 和 `thinking_signature` 顶层字段 + +```rust +#[derive(Debug, Default, Serialize, Deserialize)] +pub struct PartialMessageResponse { + pub id: Option, + pub model: Option, + pub blocks: BTreeMap, + pub block_completion: HashSet, + pub last_open_index: Option, + /// 流式累积中的部分用量信息。 + /// 使用 PartialUsage(字段为 Option)而非 Usage(字段为 u32), + /// 因为流式场景中 prompt_tokens 和 completion_tokens 可能分多次到达 + /// (如 Anthropic 的 message_delta 事件可多次下发 usage 增量)。 + pub usage: PartialUsage, + pub stop_reason: Option, + pub thinking_signature: Option, + pub is_errored: bool, + pub is_complete: bool, +} +``` + +### ContentBlockBuilder + +```rust +#[derive(Debug, Serialize, Deserialize)] +#[serde(rename_all = "snake_case")] +pub enum ContentBlockBuilder { + Text(String), + Thinking { buffer: String, signature: Option }, + Refusal(String), + /// 内聚设计:ToolUse 直接保存 arguments 字符串, + /// 而非依赖外部 PartialMessageResponse.tool_call_args HashMap。 + /// ToolCallArgumentsDelta 事件直接追加到此字段。 + ToolUse { id: String, name: String, arguments: String }, +} +``` + +### apply_to 算法 + +`StreamEvent::apply_to(&self, state: &mut PartialMessageResponse) -> bool` + +算法逻辑(基于 9c §4.4,按本次审查修订): +- `MessageStart` → 设置 id / model +- `ContentBlockStart` → 创建对应类型的 ContentBlockBuilder,按 index 分桶 +- `ContentBlockEnd` → 标记该 index 的 block 已完成 +- `TextDelta` / `ThinkingDelta` / `RefusalDelta` → 追加到最后打开的 block 缓冲区 +- `ToolCallArgumentsDelta` → 按 index 查找对应 ContentBlockBuilder::ToolUse,将 arguments 追加到其 `arguments` 字段 +- `ToolCallEnd` → 标记该 index 的 block 已完成 +- `CostUpdate` → 字段级合并 PartialUsage(只覆盖 Some 字段) +- `MessageComplete` → 设置 `is_complete = true` +- `Error` → 设置 `is_errored = true`,返回 false 终止处理 + +### finalize 算法 + +`PartialMessageResponse::finalize(self) -> Result` + +按 index 升序遍历 blocks,转换为 ContentBlock: +- Text → ContentBlock::Text +- Thinking → 如有 signature 未填充,用 `self.thinking_signature` 回填 +- Refusal → ContentBlock::Text(OpenAI refusal 合并为 Text) +- ToolUse → 从 builder 的 `arguments` 字段解析为 JSON,构造 ContentBlock::ToolUse + +最终将 `self.usage: PartialUsage` 转换为 `Usage`(缺失字段默认为 0),然后返回 MessageResponse,含 content / usage / stop_reason。 + +**关于错误处理**:如果 finalize 被调用但 `is_complete == false` 或关键字段缺失,应新增 `LlmError::Finalize(String)` 变体(在 `src/llm/error.rs` 中追加),或直接使用 `LlmError::Other(String)`。建议新增专用变体以增强错误可追溯性。 + +### 验证点 + +- `StreamEvent` 所有变体 match 穷举性 +- `PartialMessageResponse` 默认构造后状态正确 +- `apply_to(MessageStart)` → id/model 被设置 +- `apply_to` 事件序列 → `finalize()` 产生正确的 MessageResponse +- `CostUpdate` 字段级合并正确(先设置 prompt_tokens,再设置 completion_tokens,结果两者都存在) + +--- + +## 任务 5:新类型侧单元测试 + +测试按类型归属分散到对应文件: +- Message 构造 + ContentBlock roundtrip → `src/llm/types/message.rs` 的 `#[cfg(test)]` +- MessageRequest + ExtraError → `src/llm/types/request_v2.rs` 的 `#[cfg(test)]` +- MessageResponse + StreamEvent + PartialMessageResponse + JSON roundtrip → `src/llm/types/response_v2.rs` 的 `#[cfg(test)]` + +### 测试清单 + +1. **Message 构造测试** + - `Message::user_text("hello")` 产生 `Message::User { content: [ContentBlock::Text { text: "hello" }] }` + - `Message::assistant("hi")` 产生 `Message::Assistant { content: [ContentBlock::Text { text: "hi" }] }` + - `Message::system("sys")` 产生 `Message::System { content: [ContentBlock::Text { text: "sys" }] }` + - `Message::tool_result("id", "result", false)` 产生正确的 ToolResult 变体 + +2. **Message match 穷举性**(编译器验证,但显式写 match 确保不会漏变体) + +3. **PartialMessageResponse.apply_to + finalize 整合测试** + - 模拟完整的 OpenAI 流式响应(text only) + - 模拟完整的 Anthropic 流式响应(thinking + text + tool_use) + - 模拟 CostUpdate 分多次到达(字段级合并正确性) + +4. **apply_to 事件序列测试** + - MessageStart → finalize 产生正确 id/model + - ContentBlockStart(Text) → TextDelta → ContentBlockEnd → finalize 产生正确 Text block + - ContentBlockStart(ToolUse{id, name}) → ToolCallArgumentsDelta → ToolCallEnd → finalize 产生正确 ToolUse block + - ContentBlockStart(Thinking) → ThinkingDelta → ContentBlockEnd → finalize 产生正确 Thinking block + +5. **MessageComplete 事件测试** + - apply_to(MessageComplete{full_response}) 设置 `is_complete = true` + - finalize 后产生的 MessageResponse 与设置一致 + +6. **Error 事件测试** + - apply_to(Error) 设置 `is_errored = true`,返回 false + - finalize 仍然返回 Ok(即使有错误) + +7. **JSON roundtrip 测试**(对应 10 号文档 §4 任务 5 要求) + - `Message → JSON → Message`:对每个 Message 变体(System/User/UserImage/Assistant/ToolResult)分别构造实例,序列化后反序列化,验证往返不变 + - `ContentBlock → JSON → ContentBlock`:对每个 ContentBlock 变体(Text/Image/Audio/File/ToolUse/ToolResult/Thinking/Extension)分别构造实例,验证往返不变 + - `MessageRequest → JSON → MessageRequest`:构造含各种字段的完整请求,验证往返不变 + - `MessageResponse → JSON → MessageResponse`:构造含完整嵌套的响应,验证往返不变 + - `StreamEvent → JSON → StreamEvent`:对每个 event 变体验证往返不变 + +8. **ContentBlock::ToolResult 构造和序列化测试** + - 构造 `ContentBlock::ToolResult` 实例,验证字段正确性 + - 序列化后反序列化,验证 `tool_use_id`、`content`、`is_error` 不变 + +--- + +## 任务 6:修改 `src/llm/provider.rs` —— LlmProvider trait 签名切换 + +### 当前签名 + +```rust +pub trait LlmProvider: Send + Sync { + async fn chat(&self, request: ChatRequest) -> Result; + async fn chat_stream( + &self, request: ChatRequest, + ) -> Result> + Send>>, LlmError>; +} +``` + +### 目标签名 + +```rust +use crate::llm::types::request_v2::MessageRequest; +use crate::llm::types::response_v2::{StreamEvent, MessageResponse}; + +pub trait LlmProvider: Send + Sync { + async fn chat(&self, request: MessageRequest) -> Result; + async fn chat_stream( + &self, request: MessageRequest, + ) -> Result> + Send>>, LlmError>; + fn capabilities(&self) -> ProviderCapabilities; +} +``` + +**注意**: +- 新增的 `capabilities` 方法是新 trait 的一部分(参考 9c §4.2)。10 号文档 §4 实施步骤中未显式列出此方法的添加(§4 任务 6 只要求签名切换),但作为 trait 完整性要求,在 Phase 0 一并引入。此偏差已在设计评审中确认合理。 +- Phase 0 不要求实现细节——各 Provider 可以先返回默认值。 +- `ProviderCapabilities` 类型定义放在 `provider.rs` 中(与 trait 定义同文件),不分散到类型目录。 + +### ProviderCapabilities + +```rust +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct ProviderCapabilities { + pub provider_name: &'static str, + pub supported_models: Option>, + pub features: ProviderFeatures, +} + +#[derive(Debug, Clone, Default, Serialize, Deserialize)] +pub struct ProviderFeatures { + pub streaming: bool, + pub thinking: bool, + pub vision: bool, + pub audio_input: bool, + pub tool_use: bool, + pub parallel_tool_calls: bool, + pub system_prompt_in_messages: bool, + pub max_context_window: u32, +} +``` + +### import 变更 + +```rust +// 移除旧导入 +// use crate::llm::types::{ChatRequest, ChatResponse, OpenaiChatChunk}; + +// 新增新导入 +use crate::llm::types::request_v2::MessageRequest; +use crate::llm::types::response_v2::{StreamEvent, MessageResponse}; +``` + +--- + +## 任务 7:修改 `src/agent/builder.rs` —— StubProvider 适配 + +### 当前 `StubProvider` + +```rust +struct StubProvider; +#[async_trait] +impl LlmProvider for StubProvider { + async fn chat(&self, _request: ChatRequest) -> Result { + unimplemented!() + } +} +``` + +### 目标实现 + +```rust +struct StubProvider; +#[async_trait] +impl LlmProvider for StubProvider { + async fn chat(&self, _request: MessageRequest) -> Result { + unimplemented!() + } + async fn chat_stream( + &self, _request: MessageRequest, + ) -> Result> + Send>>, LlmError> { + unimplemented!() + } + fn capabilities(&self) -> ProviderCapabilities { + ProviderCapabilities { + provider_name: "stub", + supported_models: None, + features: ProviderFeatures::default(), + } + } +} +``` + +--- + +## 任务 8:修改 `src/llm/cycle.rs` —— LlmCycle 适配 + +### 8a:修改 import + +```rust +// 新类型导入 +use crate::llm::types::message::Message; +use crate::llm::types::request_v2::MessageRequest; +use crate::llm::types::response_v2::{MessageResponse, StreamEvent as NewStreamEvent}; +use crate::llm::provider::ProviderCapabilities; + +// 旧类型导入(仍用于内部消息存储和 submit_stream 处理) +use crate::llm::types::{OpenaiChatMessage, OpenaiChatChunk}; +use crate::llm::stream::StreamEvent as OldStreamEvent; + +// ⚠️ StreamEvent 命名消歧: +// - 旧 StreamEvent(src/llm/stream.rs)对应 OpenaiChatChunk 解析事件 +// - 新 StreamEvent(src/llm/types/response_v2.rs)对应高精度 IR 流式事件 +// - 本文件中,通过 as 别名区分:NewStreamEvent / OldStreamEvent +// - 实际使用 `submit_stream` 调用 provider.chat_stream 返回的是 NewStreamEvent +// - Phase 2 移除旧 StreamEvent 后,可去掉别名直接使用 StreamEvent +``` + +### 8b:修改 `build_request(&self, tools) → MessageRequest` + +将原有的 `Vec` 通过转换函数转为 `Vec`: + +```rust +fn build_request(&self, tools: &[ToolDefinition]) -> MessageRequest { + let mut ir_messages: Vec = self.messages.iter() + .map(|m| chat_message_to_message(m)) + .collect(); + + if let Some(sys_prompt) = &self.system_prompt + && !ir_messages.iter().any(|m| matches!(m, Message::System { .. })) + { + ir_messages.insert(0, Message::system_text(sys_prompt)); + } + + MessageRequest { + model: self.config.model.clone(), + messages: ir_messages, + tools: tools.to_vec(), + tool_choice: ToolChoice::Auto, + max_tokens: self.config.max_tokens, + temperature: self.config.temperature, + ..Default::default() + } +} +``` + +### 8c:转换函数 `chat_message_to_message` + +```rust +fn chat_message_to_message(msg: &OpenaiChatMessage) -> Message { + match msg { + OpenaiChatMessage::Developer { content, .. } + | OpenaiChatMessage::System { content, .. } => { + Message::System { content: content_to_blocks(content) } + } + OpenaiChatMessage::User { content, .. } => { + Message::User { content: content_to_blocks(content) } + } + OpenaiChatMessage::Assistant { content, tool_calls, .. } => { + let mut blocks = content_to_blocks(content); + // OpenAI 的 tool_calls 转为 ContentBlock::ToolUse + if let Some(calls) = tool_calls { + for call in calls { + match call { + OpenaiToolCall::Function { id, function } => { + let input: serde_json::Value = + serde_json::from_str(&function.arguments).unwrap_or_default(); + blocks.push(ContentBlock::ToolUse { + id: id.clone(), + name: function.name.clone(), + input, + }); + } + } + } + } + Message::Assistant { content: blocks } + } + OpenaiChatMessage::Tool { content, tool_call_id } => { + Message::ToolResult { + tool_call_id: tool_call_id.clone(), + content: content_to_blocks(content), + is_error: false, + } + } + // ponytail: OpenAI function_call 的 name 是函数名而非调用 ID。 + // 在 OpenAI 的实现中,旧的 function_call API 始终单次调用,name 可作为唯一标识。 + // 如需支持并行 function_call(name 重复),需在 Phase 1 Provider 实现中 + // 使用 tool_call_id 替代 name。 + OpenaiChatMessage::Function { content, name } => { + Message::ToolResult { + tool_call_id: name.clone(), + content: content_to_blocks(content), + is_error: false, + } + } + } +} +``` + +### 8d:修改 `submit()` / `submit_messages()` 返回类型 + +```rust +// 当前签名 +pub async fn submit(&mut self, prompt: String, tools: Vec) -> Result; + +// 目标签名 +pub async fn submit(&mut self, prompt: String, tools: Vec) -> Result; +``` + +**`submit_messages` 单独说明**:该方法接收 `&[OpenaiChatMessage]`(不经过 `build_request`),Phase 0 中保持输入参数类型不变。内部构造 `MessageRequest` 时,通过 `chat_message_to_message` 将传入的 `OpenaiChatMessage` 切片逐条转换为 `Vec`,然后构造 `MessageRequest { messages, ... }`。返回类型同样改为 `MessageResponse`。 + +### 8e:修改 `submit_stream()` 返回类型和逻辑 + +```rust +// 当前签名(返回 OpenaiChatChunk 流 + 旧 StreamEvent) +pub async fn submit_stream(&mut self, prompt: String, tools: Vec) + -> Result + Send>>, LlmError>; + +// 目标签名(返回新 StreamEvent 流) +// 注意:此处使用 NewStreamEvent 与 8a 中 import 别名保持一致 +pub async fn submit_stream(&mut self, prompt: String, tools: Vec) + -> Result + Send>>, LlmError>; +``` + +**内部逻辑变更**: +- 调用 `self.provider.chat_stream(request)` 直接获得 `Result` 流 +- 流处理:从消费 `OpenaiChatChunk` 改为消费 `NewStreamEvent` +- 流结束处提取 `full_response`(Phase 2 才真正使用,当前先解构 MessageComplete 拿到 MessageResponse 用于消息历史追加和 usage 统计) + +### 8f:修改 `submit_request()` 返回类型 + +```rust +// 当前 +async fn submit_request(&mut self, tools: &[ToolDefinition]) -> Result; +// 目标 +async fn submit_request(&mut self, tools: &[ToolDefinition]) -> Result; +``` + +### 8g:修改 `submit_with_tools()` 返回类型和内部逻辑 + +```rust +// 当前 +pub async fn submit_with_tools(&mut self, prompt: String, registry: &ToolRegistry) + -> Result; +// 目标 +pub async fn submit_with_tools(&mut self, prompt: String, registry: &ToolRegistry) + -> Result; +``` + +**tool 循环内部变更**: +- `has_tool_calls_in_message` 和 `extract_tool_calls_from_message` 需要适配 `MessageResponse` 类型 +- tool 结果回传时 `OpenaiChatMessage::tool_result(...)` → 但内部消息存储仍是 `Vec`(Phase 2 才切换),所以 tool 循环仍使用旧消息类型 +- MessageResponse 中的 `message` 字段是 `Message` 类型(新枚举),需要从中提取 tool_calls + +**对应调整**:当前 `submit_with_tools` 中处理的是 `response.message`(原为 `OpenaiChatMessage`),现在改为 `Message`。需要从 `Message::Assistant { content }` 中提取 `ContentBlock::ToolUse`: + +```rust +fn has_tool_calls_in_response(response: &MessageResponse) -> bool { + match &response.message { + Message::Assistant { content } => { + content.iter().any(|b| matches!(b, ContentBlock::ToolUse { .. })) + } + _ => false, + } +} + +fn extract_tool_calls_from_response(response: &MessageResponse) -> Vec<(String, String, String)> { + match &response.message { + Message::Assistant { content } => { + content.iter().filter_map(|b| { + if let ContentBlock::ToolUse { id, name, input } = b { + Some((id.clone(), name.clone(), serde_json::to_string(input).unwrap_or_default())) + } else { + None + } + }).collect() + } + _ => vec![], + } +} +``` + +### 8h:修改 `submit_with_tools()` 中 tool 结果回传 + +当前使用 `OpenaiChatMessage::tool_result(...)`——因为内部 `self.messages` 仍是 `Vec`,所以保持使用旧类型。MessageResponse 的 tool_calls 提取改用 8g 中的新函数。 + +### 8i:修改 `build_request` 中 hooks context 的类型 + +`HookContext::with_request` 目前接受 `&ChatRequest`。改为接受 `&MessageRequest`。相应修改 `src/llm/hooks.rs` 中的 `HookContext` 类型定义。 + +### 8j:修改 `submit_stream()` 后处理 + +流结束后: +- 监听 `MessageComplete` 事件提取 `full_response: MessageResponse` +- 将 `full_response.message` 转为 `OpenaiChatMessage`(新→旧)推入 `self.messages` +- 累加 usage + +这是 Phase 0 的临时逻辑——Phase 2 会移除这个转换。 + +### 8k:转换函数 `message_to_chat_message` + +在 Phase 0 中需要将新 `Message` 转回 `OpenaiChatMessage`(用于内部消息存储): + +```rust +fn message_to_chat_message(msg: &Message) -> OpenaiChatMessage { + match msg { + Message::System { content } => OpenaiChatMessage::System { + content: blocks_to_content_field(content), + name: None, + }, + Message::User { content } => OpenaiChatMessage::User { + content: blocks_to_content_field(content), + name: None, + }, + // TODO(Phase2): UserImage → User 消息 + Image content block 的真实映射。 + // Phase 0 中内部消息存储仍为 Vec,不支持图片 content block, + // 因此 UserImage 转换为空消息。此阶段任何涉及 UserImage 回流的测试/功能 + // 将丢失图片数据。Phase 2 切换为 Vec 后此函数整体移除。 + Message::UserImage { data: _, mime_type: _, detail: _ } => { + OpenaiChatMessage::User { + content: ContentField::Array(vec![]), + name: None, + } + } + Message::Assistant { content } => { + // 从 content 中提取 tool_use 块 + let mut text_blocks = vec![]; + let mut tool_call_blocks = vec![]; + for block in content { + match block { + ContentBlock::Text { text } => text_blocks.push(text.clone()), + ContentBlock::ToolUse { id, name, input } => { + tool_call_blocks.push(OpenaiToolCall::Function { + id: id.clone(), + function: FunctionCall { + name: name.clone(), + arguments: serde_json::to_string(input).unwrap_or_default(), + }, + }); + } + // ponytail: 静默跳过 thinking/refusal/image/audio/file 等非 OpenAI 原生 block。 + // Phase 0 内部消息存储仍使用 Vec,这些 block 在回传时被截断。 + // 影响范围:多轮对话中,thinking block 不会出现在后续请求的上下文中。 + // Phase 2 切换为 Vec 后此函数整体移除,自然解决。 + _ => {} + } + } + let text = text_blocks.join(""); + OpenaiChatMessage::Assistant { + content: if text.is_empty() { + ContentField::Array(vec![]) + } else { + ContentField::String(text) + }, + refusal: None, + name: None, + tool_calls: if tool_call_blocks.is_empty() { None } else { Some(tool_call_blocks) }, + } + } + Message::ToolResult { tool_call_id, content, is_error: _ } => { + OpenaiChatMessage::Tool { + content: blocks_to_content_field(content), + tool_call_id: tool_call_id.clone(), + } + } + } +} +``` + +--- + +## 任务 9:编译驱动适配上游 + +### 9a:`src/agent/session.rs` + +- `submit_turn()` 返回类型从 `Result` 改为 `Result` +- `response.usage` 引用适配(`MessageResponse.usage` 字段与 `ChatResponse.usage` 类型同为 `Usage`,无需修改) +- `response.message` 从 `OpenaiChatMessage` 变为 `Message`——测试代码中 `match response.message` 的分支需要适配 +- 测试模块(`src/agent/session.rs` 中的 `#[cfg(test)]`)中的 `MockProvider` 需要适配新 trait 签名(返回 `MessageResponse` 而非 `ChatResponse`),并新增 `chat_stream` 和 `capabilities` 方法实现 +- `assistant_text()` 辅助函数改为构造 `MessageResponse` 而非 `ChatResponse` + +### 9b:`src/agent/runtime.rs` + +- 无直接修改(只引用 `Arc`,trait 定义变更不影响持有者) +- 但需验证 `RuntimeBundle` 中的 `provider: Arc` 在 trait 签名变更后编译通过 + +### 9c:`src/agent/error.rs` + +- 无预期变更(`AgentError` 不直接引用 `ChatResponse` 或相关类型) + +### 9d:`src/memory/conversation.rs` + +- `messages: Vec` 保持不动(Phase 2 才切换为 `Vec`) +- 但 `ConversationMemory` 的 `add_message()` 和 `get_history()` 方法签名需要确认是否需要适配 +- `serde_json::from_str::` 反序列化保持使用旧类型 + +### 9e:`src/prompt/composer.rs` + +- `build_request()` 返回 `OpenaiChatRequest` 保持不变(旧类型仍存在) +- `validate_messages()` 继续使用 `&[OpenaiChatMessage]` 保持不变 + +### 9f:`src/llm/cycle/retry.rs` + +- 检查是否有对 `ChatRequest` / `ChatResponse` 的引用,如有则更新为 `MessageRequest` / `MessageResponse` +- 具体检查点:`retry_with_backoff` 函数的请求参数类型、`should_retry` 函数的响应/错误类型、以及 `retry_policy` 配置中涉及的类型引用 + +### 9g:`src/llm/hooks.rs` + +- `HookContext::with_request` 参数类型从 `&ChatRequest` 改为 `&MessageRequest` +- **决策锁定**:直接修改 `with_request` 签名,不做双方法兼容。理由:Phase 0 完成后旧 `ChatRequest` 只在旧 Provider 内部使用(`src/llm/provider/openai.rs`),hook 层不应引用旧 request 类型。 +- 同步更新 `HookContext` 中所有引用 `ChatRequest` 的字段和方法 + +### 9h:`src/llm/compact.rs` + +- 暂时不修改(`estimate_message_tokens` / `microcompact` 继续使用 `OpenaiChatMessage`) +- Phase 2 才做切换 + +### 9i:`examples/` 目录 + +- 运行 `grep -r "ChatResponse\|ChatRequest" examples/` 预检 +- 如现有 example 直接引用旧类型,同步适配为 `MessageResponse` / `MessageRequest` +- 如果 example 中构造了旧类型实例(`OpenaiChatMessage` 等),由于旧类型本身仍存在,只需更新引用新类型的部分 + +--- + +## 验证方式 + +1. **编译检查**:`cargo build` 通过,无类型错误 +2. **单元测试**:`cargo test` 全部通过 +3. **clippy 检查**:`cargo clippy` 无新增警告 +4. **diff 确认**:`git diff --stat` 确认修改范围符合预期 + +## 回滚方案 + +| 触发条件 | 操作 | +|---------|------| +| 新类型设计发现重大缺陷(如 Message 枚举扁平度不足) | 回退 git 到 Phase 0 开始前的 commit,保留 9 系文档作为参照,重启设计评审 | +| 临时转换层(`chat_message_to_message` / `message_to_chat_message`)逻辑错误 | 修复转换函数逻辑,不修改类型定义 | +| 单个 Provider 的适配不合理 | 将该 Provider 回退为 `unimplemented!()`(当前状态),不影响其他组件 | +| Phase 0 完成时发现未预见的设计问题 | **不打乱整体进度**:在不动已有文件的前提下,直接原地修改新类型文件重新迭代,不需要整个回退到 Phase 0 之前。Phase 0 结束时打 tag `types-v2-prototype` 作为 checkpoint | + +**风险储备**: +- `ChatRequest` / `ChatResponse` 等旧类型保留不动,不主动删除。Phase 2 切换后通过 `#[deprecated]` 标记引导迁移 +- 如果 `OpenaiProvider` 保留旧实现,旧 trait 方法名不变、只是签名变,不影响共存 + +## 开放事项 + +以下事项已在 10 号文档中讨论,Phase 0 实施前需确认: + +- [x] `CompletionTokensDetails` 和 `PromptTokensDetails` — 已确认存在于 `src/llm/types/usage.rs:14-32`,通过 `types/mod.rs:21` 重导出,无需新增 +- [ ] Phase 0 完成后确认 `git diff --stat` 只涉及预期文件,无意外修改 +- [ ] `types/mod.rs` 中旧 `pub type Message = OpenaiChatMessage` 的引用面——执行 `grep "types::Message" src/ -r` 确认是否无处引用(实施前基线确认中已包含此步骤) diff --git a/docs/10b-phase1-provider-adaptation.md b/docs/10b-phase1-provider-adaptation.md new file mode 100644 index 0000000..1d5cf45 --- /dev/null +++ b/docs/10b-phase1-provider-adaptation.md @@ -0,0 +1,1133 @@ +# Phase 1 实施计划:Provider 适配 + +> **所属方案**:[10-llm-provider-refinement.md](10-llm-provider-refinement.md) +> +> **前置条件**:Phase 0 已完成,`LlmProvider` trait 签名已切换为 `chat(MessageRequest) → MessageResponse` +> +> **产出依赖**:Phase 2(LlmCycle 简化)依赖本阶段完成的 Provider 实现 + +## Phase 0 就绪检测 + +**前置条件**:Phase 0 引入的 `MessageRequest`/`MessageResponse`/`StreamEvent` 等新类型必须已在代码库中落地。 +进入 Phase 1 执行前,必须满足以下全部条件: + +| 检测项 | 检查方式 | 预期结果 | +|--------|---------|---------| +| `request_v2.rs` / `response_v2.rs` 存在 | `ls src/llm/types/` | 文件存在 | +| `pub type Message = OpenaiChatMessage` 别名已移除 | `grep "pub type Message" src/llm/types/mod.rs` | 无匹配 | +| `pub type ChatRequest = OpenaiChatRequest` 别名已移除 | `grep "pub type ChatRequest" src/llm/types/mod.rs` | 无匹配 | +| `StreamEvent` 包含 `ContentBlockStart` 新变体 | `grep "ContentBlockStart" src/llm/stream.rs` | 有匹配(Phase 0 引入) | +| `ContentBlock` 是新类型的泛化 enum(非 `OpenaiContentPart` alias) | `grep "pub enum ContentBlock" src/llm/types/message.rs` | 有匹配(Phase 0 引入) | +| `LlmProvider` trait 使用新签名 | `grep "async fn chat" src/llm/provider.rs` | 签名为 `fn chat(&self, request: MessageRequest) -> Result` | +| `cargo build` 编译通过 | `cargo build 2>&1` | 无错误 | + +**以上条件未满足时,禁止启动 Phase 1 开发。** + +## Phase 0 未完成时的降级执行路径 + +如果 Phase 0 因各种原因延迟,以下工作可独立于 Phase 0 先行完成(**不依赖新类型**): + +| 可独立执行的内容 | 不依赖 Phase 0 的原因 | 过渡策略 | +|-----------------|---------------------|---------| +| `ProviderType` enum 扩展(+Anthropic/OpenaiResponse) | enum 扩展不涉及 trait 签名变更 | 先加变体,`create_provider` 返回 `Err(LlmError::Other(...))` | +| AnthropicProvider 的 HTTP 调用层 | 使用 `reqwest::Client` 原始调用,不依赖 `LlmProvider` trait | 独立函数(`async fn anthropic_chat(...)`),Phase 0 完成后包装为 trait impl | +| SSE 字节解析通用层(`SseByteStream`) | 纯字节处理,无关类型系统 | 独立模块,Phase 0 完成后挂接到 `StreamEvent` 映射 | +| `wiremock` dev-dependency 添加 | 仅 Cargo.toml 变更 | 先加依赖,测试代码在 Phase 0 后补齐 | + +当 Phase 0 完成后,上述先行部分直接集成到 Phase 1 主流程,无需重写。 + +--- + +## 目标 + +重写 `OpenaiProvider`(使用新类型),新增 `AnthropicProvider`。DeepSeek/Qwen 作为 OpenAI-compatible 协议实现一并纳入。本 Phase 结束时,每个 Provider 都有 mock 覆盖的 chat + chat_stream 基本路径测试。 + +## 设计决策 + +### ProviderType enum 变更 + +当前 enum(3 个变体)→ 目标 enum(5 个变体): + +```rust +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub enum ProviderType { + OpenaiChat, // 原 OpenAI(标准 /v1/chat/completions) + OpenaiResponse, // OpenAI Response API(新增,独特协议) + Anthropic, // Anthropic Messages API(新增) + DeepSeek, // DeepSeek(OpenAI-compatible) + Qwen, // 通义千问(OpenAI-compatible) +} +``` + +**名称变更注意**:当前 `OpenAI` 改名为 `OpenaiChat`,精确区分 Chat API 和 Response API。 + +### OpenAI-compatible 复用策略 + +**推荐方案:参数化 GenericOpenaiProvider** + +```rust +pub struct GenericOpenaiProvider { + http_client: Client, + base_url: String, + api_key: String, + model: String, + provider_name: &'static str, + /// 额外请求头,由 Provider newtype 构造器传入(如 Qwen 的 X-DashScope-SSE: enable) + extra_headers: Vec<(String, String)>, +} +``` + +构造器定义(供 DeepSeek/Qwen newtype 调用): +```rust +impl GenericOpenaiProvider { + pub fn new_with_name( + base_url: String, api_key: String, model: String, + provider_name: &'static str, + ) -> Self { /* ... */ } + + pub fn new_with_name_and_headers( + base_url: String, api_key: String, model: String, + provider_name: &'static str, + extra_headers: Vec<(String, String)>, + ) -> Self { /* ... */ } +} +``` + +DeepSeek 和 Qwen 共用同一实现,仅配置不同: +- `GenericOpenaiProvider::new_with_name("https://api.deepseek.com", api_key, model, "deepseek")` +- `GenericOpenaiProvider::new_with_name("https://dashscope.aliyuncs.com/compatible-mode/v1", api_key, model, "qwen")` + +**extra_headers 注入**:在 `chat()` / `chat_stream()` 的 HTTP 请求构造中,遍历 `self.extra_headers` 注入到 `reqwest::RequestBuilder`: +```rust +let mut request_builder = self.http_client.post(&url) + .header("Authorization", format!("Bearer {}", self.api_key)); +for (key, value) in &self.extra_headers { + request_builder = request_builder.header(key.as_str(), value.as_str()); +} +``` + +**备选方案**:如果不想引入 GenericOpenaiProvider 的抽象层,OpenaiChatProvider 直接覆盖所有 OpenAI-compatible 场景,DeepSeek/Qwen 的 Provider 只是调用 OpenaiChatProvider 的不同实例。 + +**实施时决定**:先实现 `OpenaiChatProvider`,确认它足够的通用性后,将 DeepSeek/Qwen 简化为配置别名。 + +### OpenaiResponseProvider 范围 + +`OpenaiResponseProvider` 在本 Phase 只覆盖**核心对话能力**(create response + 流式 + 工具调用),内置工具(web_search, file_search)、`previous_response_id` 续写、`store` 等特性通过 `MessageRequest.extra` 传递。如果资源有限,可延迟到 Phase 2 之后开发,不影响其他 Provider。 + +### HTTP mock 策略 + +- 新增 dev-dependency: [`wiremock`](https://crates.io/crates/wiremock) +- 每个 Provider 的测试模块中,用 `MockServer` 启动 mock 服务端 +- `OpenaiChatProvider` mock 端点:`/chat/completions`(SSE 流和 JSON 响应) +- `AnthropicProvider` mock 端点:`/v1/messages`(SSE 事件序列) +- Provider 构造时 `base_url` 指向 `mock_server.uri()` +- 测试不依赖真实网络 + +--- + +## 涉及文件 + +| 操作 | 文件 | 说明 | +|------|------|------| +| 修改 | `src/llm/provider.rs` | 扩展 `ProviderType` enum、更新 `create_provider`、更新 `FromStr` | +| 修改 | `src/llm/provider/openai.rs` | 重写:接受 MessageRequest,返回 MessageResponse/StreamEvent;新增 `GenericOpenaiProvider` | +| 新增 | `src/llm/provider/openai_compat.rs` | DeepSeek/Qwen 的 newtype 包装 + delegate macro(替代独立 deepseek.rs/qwen.rs) | +| 新增 | `src/llm/provider/anthropic.rs` | Anthropic Messages API 实现(含独立错误映射) | +| **新增** | `src/llm/convert.rs` | 公共转换模块:`from_openai()` / `to_openai()` / `content_to_blocks()` / `blocks_to_content()` | +| 修改 | `src/llm/provider/registry.rs` | 适配新的 ProviderType 和 trait 签名 | +| 修改 | `Cargo.toml` | 新增 wiremock dev-dependency | + +--- + +## 任务 1:扩展 ProviderType enum 和工厂函数 + +### 1a:扩展 enum + +```rust +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub enum ProviderType { + OpenaiChat, + OpenaiResponse, + Anthropic, + DeepSeek, + Qwen, +} +``` + +### 1b:更新 FromStr + +```rust +impl std::str::FromStr for ProviderType { + type Err = String; + fn from_str(s: &str) -> Result { + match s.to_lowercase().as_str() { + "openai" | "openaichat" => Ok(ProviderType::OpenaiChat), + "openai-response" | "response" => Ok(ProviderType::OpenaiResponse), + "anthropic" | "claude" => Ok(ProviderType::Anthropic), + "deepseek" => Ok(ProviderType::DeepSeek), + "qwen" | "dashscope" | "tongyi" => Ok(ProviderType::Qwen), + _ => Err(format!("未知的 Provider 类型: {}", s)), + } + } +} +``` + +### 1c:更新 `create_provider` 工厂函数 + +```rust +pub fn create_provider( + provider_type: ProviderType, + config: ProviderConfig, +) -> Result, LlmError> { + match provider_type { + ProviderType::OpenaiChat => Ok(Box::new(openai::OpenaiChatProvider::new( + config.base_url, config.api_key, config.model, + ))), + ProviderType::OpenaiResponse => { + // 注意:此处不返回 unimplemented!(),避免运行时 panic。 + // 返回 Err 让调用方可以优雅地降级。 + Err(LlmError::Other("OpenaiResponse Provider 尚未实现,请使用 OpenaiChat".into())) + } + ProviderType::Anthropic => Ok(Box::new(anthropic::AnthropicProvider::new( + config.base_url, config.api_key, config.model, + ))), + ProviderType::DeepSeek => Ok(Box::new(openai_compat::DeepSeekProvider::new( + config.base_url, config.api_key, config.model, + ))), + ProviderType::Qwen => Ok(Box::new(openai_compat::QwenProvider::new( + config.base_url, config.api_key, config.model, + ))), + } +} +``` + +--- + +## 任务 2:重写 OpenaiChatProvider + +### 2a:结构体定义 + +```rust +pub struct OpenaiChatProvider { + http_client: Client, + base_url: String, + api_key: String, + model: String, +} +``` + +### 2b:非流式 chat() + +```rust +async fn chat(&self, request: MessageRequest) -> Result { + let url = format!("{}/chat/completions", self.base_url.trim_end_matches('/')); + + // 1. MessageRequest → OpenaiChatRequest + let openai_req = self.convert_request(request)?; + + // 2. HTTP POST + let response = self.http_client + .post(&url) + .header("Authorization", format!("Bearer {}", self.api_key)) + .json(&openai_req) + .send() + .await + .map_err(Self::map_reqwest_error)?; + + // 3. 错误处理(复用现有逻辑) + let status = response.status(); + if !status.is_success() { + return Err(self.handle_error(response).await); + } + + // 4. 解析 OpenaiChatResponse + let body_text = response.text().await.unwrap_or_default(); + let chat_response: OpenaiChatResponse = serde_json::from_str(&body_text) + .map_err(|e| LlmError::Other(format!("响应解析失败: {}", e)))?; + + // 5. OpenaiChatResponse → MessageResponse + self.convert_response(chat_response) +} +``` + +### 2c:流式 chat_stream() + +```rust +async fn chat_stream( + &self, request: MessageRequest, +) -> Result> + Send>>, LlmError> { + // 1. MessageRequest → OpenaiChatRequest + 设置 stream=true + let mut openai_req = self.convert_request(request)?; + openai_req.stream = Some(true); + openai_req.stream_options = Some(StreamOptions { + include_usage: Some(true), + include_obfuscation: None, + }); + + // 2. HTTP POST 获取 SSE 流 + let url = format!("{}/chat/completions", self.base_url.trim_end_matches('/')); + let response = self.http_client + .post(&url) + .header("Authorization", format!("Bearer {}", self.api_key)) + .json(&openai_req) + .send() + .await + .map_err(Self::map_reqwest_error)?; + + let status = response.status(); + if !status.is_success() { + return Err(self.handle_error(response).await); + } + + // 3. 将 byte stream 转为 StreamEvent 流 + let byte_stream = response.bytes_stream().map(|r| { + r.map_err(|e| LlmError::Other(format!("流式读取失败: {}", e))) + }); + + Ok(Box::pin(OpenaiChunkToEventStream::new(byte_stream))) +} +``` + +### 2d:转换逻辑 `convert_request()` + +`MessageRequest` → `OpenaiChatRequest`: + +```rust +fn convert_request(&self, request: MessageRequest) -> Result { + // messages: Vec → Vec + let messages: Vec = request.messages + .iter() + .map(|m| self.message_to_chat_message(m)) + .collect(); + + // 构造 OpenaiChatRequest + Ok(OpenaiChatRequest { + model: request.model, + messages, + max_tokens: request.max_tokens, + temperature: request.temperature, + top_p: request.top_p, + tools: if request.tools.is_empty() { None } else { + Some(request.tools.into_iter().map(|t| OpenaiTool::Function { + function: t, + }).collect()) + }, + // 假设:MessageRequest.tool_choice 为 ToolChoice(非 Option), + // 与 OpenaiChatRequest.tool_choice: Option 自然匹配。 + // 如果 Phase 0 定义为 Option,去掉 Some(...) 包装。 + tool_choice: Some(request.tool_choice), + stream: None, // 由调用方在 chat_stream 中设置 + stop: Some(StopSequence::Multiple(request.stop_sequences)), + // 从 extra 读取 OpenAI 特有参数 + frequency_penalty: request.get_extra_opt("frequency_penalty"), + presence_penalty: request.get_extra_opt("presence_penalty"), + seed: request.get_extra_opt("seed"), + response_format: request.get_extra_opt("response_format"), + parallel_tool_calls: request.get_extra_opt("parallel_tool_calls"), + // o-series 模型使用 max_completion_tokens 而非 max_tokens + // 当 request.extra["max_completion_tokens"] 存在时覆盖 max_tokens + max_completion_tokens: request.get_extra_opt("max_completion_tokens"), + ..Default::default() + }) +} + +**Extra 字段清单**(从 `MessageRequest.extra` 透传到 `OpenaiChatRequest`): + +| Extra Key | 目标字段 | 说明 | +|-----------|---------|------| +| `frequency_penalty` | `frequency_penalty` | 频率惩罚 | +| `presence_penalty` | `presence_penalty` | 存在惩罚 | +| `seed` | `seed` | 随机种子 | +| `response_format` | `response_format` | JSON mode / Structured Outputs(企业高频需求) | +| `parallel_tool_calls` | `parallel_tool_calls` | 是否允许并行工具调用(默认 true) | +| `max_completion_tokens` | `max_completion_tokens` | o-series 模型专用,覆盖 `max_tokens` | +| `openai_organization` | HTTP header `OpenAI-Organization` | 企业组织标识 | +| `openai_project` | HTTP header `OpenAI-Project` | 企业项目标识 | + +HTTP header 参数在 `chat()`/`chat_stream()` 中从 extra 提取后注入请求头。 + +### 公共转换模块 + +`Message` ↔ `OpenaiChatMessage` 的转换函数在 Phase 0 的 `LlmCycle` 和 Phase 1 的 `OpenaiChatProvider` 中都会用到。 +为避免逻辑漂移,应抽取为公共模块: + +```rust +// src/llm/convert.rs — 跨 Provider 类型转换模块 + +/// OpenaiChatMessage → Message +pub fn from_openai(msg: &OpenaiChatMessage) -> Message { ... } + +/// Message → OpenaiChatMessage +pub fn to_openai(msg: &Message) -> OpenaiChatMessage { ... } + +/// ContentField → Vec +pub fn content_to_blocks(field: &ContentField) -> Vec { ... } + +/// Vec → ContentField +pub fn blocks_to_content(blocks: &[ContentBlock]) -> ContentField { ... } +``` + +Phase 0 的 `cycle.rs` 和 Phase 1 的所有 Provider 统一引用此模块。 + +### 2e:转换逻辑 `convert_response()` + +`OpenaiChatResponse` → `MessageResponse`: + +```rust +fn convert_response(&self, response: OpenaiChatResponse) -> Result { + let choice = response.choices.into_iter().next() + .ok_or_else(|| LlmError::Other("响应中没有 choices".to_string()))?; + + let message = self.chat_message_to_message(&choice.message); + let stop_reason = match choice.finish_reason { + Some(FinishReason::Stop) => StopReason::Stop, + Some(FinishReason::Length) => StopReason::Length, + Some(FinishReason::ToolCalls) => StopReason::ToolUse, + Some(FinishReason::ContentFilter) => StopReason::ContentFilter, + _ => StopReason::Other, + }; + + Ok(MessageResponse { + id: response.id, + model: response.model, + message, + usage: response.usage, + stop_reason, + extra: HashMap::new(), + }) +} +``` + +### 2f:转换逻辑 `chat_message_to_message()` + +`OpenaiChatMessage` → `Message`(与 Phase 0 中定义的 `chat_message_to_message` 一致,可复用): + +```rust +fn chat_message_to_message(&self, msg: &OpenaiChatMessage) -> Message { + // 与 Phase 0 中定义的转换函数一致 + match msg { + OpenaiChatMessage::System { content, .. } + | OpenaiChatMessage::Developer { content, .. } => { + Message::System { content: content_to_blocks(content) } + } + OpenaiChatMessage::User { content, .. } => { + Message::User { content: content_to_blocks(content) } + } + OpenaiChatMessage::Assistant { content, tool_calls, .. } => { + let mut blocks = content_to_blocks(content); + if let Some(calls) = tool_calls { + for call in calls { + if let OpenaiToolCall::Function { id, function } = call { + let input: Value = serde_json::from_str(&function.arguments).unwrap_or_default(); + blocks.push(ContentBlock::ToolUse { + id: id.clone(), + name: function.name.clone(), + input, + }); + } + } + } + Message::Assistant { content: blocks } + } + OpenaiChatMessage::Tool { content, tool_call_id } => { + Message::ToolResult { + tool_call_id: tool_call_id.clone(), + content: content_to_blocks(content), + is_error: false, + } + } + OpenaiChatMessage::Function { content, name } => { + Message::ToolResult { + tool_call_id: name.clone(), + content: content_to_blocks(content), + is_error: false, + } + } + } +} +``` + +### 2g:转换逻辑 `message_to_chat_message()` + +`Message` → `OpenaiChatMessage`(与 Phase 0 一致): + +```rust +fn message_to_chat_message(&self, msg: &Message) -> OpenaiChatMessage { + // 与 Phase 0 中定义的 message_to_chat_message 一致 + // 详见 Phase 0 任务 8k +} +``` + +### 2h:流式 SSE → StreamEvent 转换 + +`OpenaiChunkToEventStream` 的流程: +1. 读取字节流,按 `\n` 分割 SSE 行 +2. 解析 `data:` 前缀的 JSON 为 `OpenaiChatChunk` +3. 将 `OpenaiChatChunk` 转换为 `StreamEvent` + +#### 状态机设计 + +因为 OpenAI 流式没有显式 block 边界,需要内部状态机追踪当前活跃的 block: + +``` +状态: Idle | InTextBlock { index } | InToolCall { index, tool_index } | InRefusalBlock { index } + +事件处理: + Idle + choices[0].delta.role == "assistant" + → 发出 MessageStart { id, model },保持 Idle + + Idle/InTextBlock + delta.content 首次出现 + → 发出 ContentBlockStart { index: next_block_index, block_type: Text } + → 发出 TextDelta { text } + → 状态: InTextBlock { index } + + InTextBlock + delta.content 后续值 + → 发出 TextDelta { text } + + InTextBlock + finish_reason 出现(且无 tool_calls) + → 发出 ContentBlockEnd { index } + → 状态: Idle + + Idle + delta.refusal 首次出现 + → 发出 ContentBlockStart { index: next_block_index, block_type: Refusal } + → 发出 RefusalDelta { text } + → 状态: InRefusalBlock { index } + + InRefusalBlock + delta.refusal 后续值 + → 发出 RefusalDelta { text } + + InRefusalBlock + delta.content 或 delta.tool_calls + → 忽略(refusal block 内部不应有其他内容类型;如确有,按防御性编程忽略) + + InRefusalBlock + finish_reason 出现 + → 发出 ContentBlockEnd { index } + → 状态: Idle + + Idle/InTextBlock/InToolCall/InRefusalBlock + delta.tool_calls[i] 首次出现 + → 发出 ContentBlockStart { index: next_block_index, block_type: ToolUse { id, name } } + → 发出 ToolCallArgumentsDelta { index, arguments } + → 状态: InToolCall { index, tool_index: i } + + InToolCall + delta.tool_calls[i].arguments 后续值 + → 发出 ToolCallArgumentsDelta { index, arguments } + + InToolCall + delta.tool_calls[i] 结束(arguments 为空且 finish_reason 出现) + → 发出 ToolCallEnd { index } + → 状态: Idle + + 任何状态 + finish_reason 出现(且 tool_calls 已结束) + → 如有未结束的 TextBlock/RefusalBlock:发出 ContentBlockEnd { index } + → 状态: Idle + + 任何状态 + usage 有值 + → 发出 CostUpdate { usage } + + 收到 data: [DONE] + → 发出 MessageComplete { full_response: partial.finalize()? } +``` + +#### index 管理规则 + +| 规则 | 说明 | +|------|------| +| text block 的 index | `next_block_index` 从 0 开始递增,每个新 text block +1 | +| tool_call 的 index | **使用 chunk 中 `choices[0].delta.tool_calls[i].index` 字段**,而非 Provider 自增计数器 | +| 并行 tool_calls | 多个 tool_call 可能同时出现在同一个 chunk 中,各自有独立 index | +| index 冲突 | text block 和 tool_call 共用一套 index 空间(OpenAI 不混合同一切片) | + +#### 关键边界场景 + +| 场景 | 处理方式 | +|------|---------| +| 同一切片同时出现 `content` + `tool_calls` | 先处理 `content` 的 TextDelta(或在 InTextBlock 中追加),再处理 `tool_calls` | +| `refusal` 字段 | 第一次出现时:发出 `ContentBlockStart { block_type: Refusal }` + `RefusalDelta { text }` | +| 仅 tool_calls 无内容 | skip ContentBlockStart(Text) 直接进入 tool_call 处理 | +| choices 数组为空的中间 chunk | 忽略,仅触发 usage 或 finish_reason 解析 | + +**MessageComplete 事件的构造**: +使用 `PartialMessageResponse` + `apply_to` 组装,然后 `finalize()` 得出 `MessageResponse`。当收到最后一个 chunk 或 `[DONE]` 信号时,调用 `partial.finalize()` 并发出 `MessageComplete { full_response }`。 + +--- + +## 任务 3:实现 AnthropicProvider + +### 3a:配置 + +```rust +pub struct AnthropicProvider { + http_client: Client, + base_url: String, // 默认: https://api.anthropic.com + api_key: String, + model: String, // 默认: claude-sonnet-4-20250514 + anthropic_version: String, // 默认: 2023-06-01 +} + +impl AnthropicProvider { + pub fn new(base_url: String, api_key: String, model: String) -> Self { + // 使用 expect 提供有意义的错误信息,避免裸 unwrap 导致难以追踪的 panic + let key_header = HeaderValue::from_str(&api_key) + .expect("Anthropic API key 包含无效的 HTTP 头部字符(如控制字符)"); + let version_header = HeaderValue::from_static("2023-06-01"); + + let http_client = Client::builder() + .timeout(Duration::from_secs(120)) + .default_headers({ + let mut headers = HeaderMap::new(); + headers.insert("x-api-key", key_header); + headers.insert("anthropic-version", version_header); + headers + }) + .build() + .expect("创建 HTTP 客户端失败"); + + Self { + http_client, + base_url: if base_url.is_empty() { + "https://api.anthropic.com".to_string() + } else { + base_url + }, + api_key, + model, + anthropic_version: "2023-06-01".to_string(), + } + } +} +``` + +#### max_tokens 必填兜底 + +Anthropic Messages API **要求 `max_tokens` 为必填字段**,而 `MessageRequest.max_tokens` 的类型是 `Option`。 +当用户未设置时,Provider 必须提供一个合理的默认值: + +```rust +fn resolve_max_tokens(&self, request: &MessageRequest) -> u32 { + request.max_tokens.unwrap_or(4096) // Anthropic 推荐的合理默认值 +} +``` + +默认值选择:`4096` 是 Anthropic 官方推荐的无特殊需求时的合理值,覆盖大部分场景。 +如果希望更精确,也可以从 `capabilities().features.max_context_window / 10` 计算(约 20000), +但对于短对话 20000 可能偏高,`4096` 更安全。 + +### 3b:非流式 chat() — MessageRequest → Anthropic Messages API 请求体 + +Anthropic Messages API 请求格式: +```json +{ + "model": "claude-sonnet-4-20250514", + "max_tokens": 1024, + "system": [{"type": "text", "text": "..."}], + "messages": [ + {"role": "user", "content": [{"type": "text", "text": "hello"}]}, + {"role": "assistant", "content": [{"type": "text", "text": "hi"}]} + ], + "tools": [...] +} +``` + +**转换要点**: +- `Message::System` → 提取为请求体的 `system` 参数(Anthropic 中 system 不在 messages 中) +- `Message::User { content }` → `{ role: "user", content }`(ContentBlock 直接映射) +- `Message::UserImage { data, mime_type, detail }` → `{ role: "user", content: [{ type: "image", source: { type: "base64", media_type: mime_type, data } }] }` +- `Message::Assistant { content }` → `{ role: "assistant", content }`(ContentBlock 直接映射) +- `Message::ToolResult { tool_call_id, content, is_error }` → `{ role: "user", content: [{ type: "tool_result", tool_use_id: tool_call_id, content, is_error }] }` +- `ThinkingConfig` → 请求体增加 `thinking: { type: "enabled", budget_tokens: N }` + +### 3c:Anthropic 响应 → MessageResponse + +Anthropic Messages API 响应格式: +```json +{ + "id": "msg_01...", + "type": "message", + "role": "assistant", + "content": [ + {"type": "text", "text": "..."}, + {"type": "tool_use", "id": "toolu_01...", "name": "get_weather", "input": {...}}, + {"type": "thinking", "thinking": "..."} + ], + "model": "claude-sonnet-4-20250514", + "stop_reason": "end_turn", + "usage": { + "input_tokens": 100, + "output_tokens": 50 + } +} +``` + +**转换要点**: +- `stop_reason`: `end_turn` → `Stop`, `max_tokens` → `MaxTokens`, `tool_use` → `ToolUse`, `stop_sequence` → `StopSequence` +- `content` 中的 block 直接映射为 `ContentBlock`(text → Text, tool_use → ToolUse, thinking → Thinking) +- `usage`: `input_tokens` → `prompt_tokens`, `output_tokens` → `completion_tokens` + +### 3g:错误映射 + +Anthropic API 的错误响应格式与 OpenAI 不同,需要独立的错误解析逻辑: + +**Anthropic 错误响应格式**: +```json +{ + "type": "error", + "error": { + "type": "authentication_error", + "message": "Invalid API key" + } +} +``` + +**错误类型映射**: + +| HTTP 状态码 | error.type | 映射到 LlmError | +|-------------|-----------|----------------| +| 400 | `invalid_request_error` | `LlmError::Request { status: 400, body }` | +| 401 | `authentication_error` | `LlmError::Authentication(body)` | +| 403 | `permission_error` | `LlmError::Authentication(body)` | +| 404 | `not_found_error` | `LlmError::Request { status: 404, body }` | +| 429 | `rate_limit_error` | `LlmError::RateLimit { retry_after }` | +| 500 | `api_error` | `LlmError::Request { status: 500, body }` | +| 529 | `overload_error` | `LlmError::RateLimit { retry_after: Some(Duration::from_secs(30)) }` | + +**529 特殊处理**:Anthropic 在过载时返回 529,应当映射到 `RateLimit` 而不是 `>=500` 的 `Request`, +这样上游的 retry 机制可以正确处理自动重试。 + +```rust +async fn handle_anthropic_error(&self, response: Response) -> LlmError { + let status = response.status().as_u16(); + let body = response.text().await.unwrap_or_default(); + + match status { + 401 | 403 => LlmError::Authentication(body), + 429 | 529 => { + let retry_after = response + .headers() + .get("retry-after") + .and_then(|v| v.to_str().ok()) + .and_then(|v| v.parse::().ok()) + .map(Duration::from_secs); + LlmError::RateLimit { retry_after } + } + _ => LlmError::Request { status, body }, + } +} +``` + +### 3h:流式 chat_stream() + +Anthropic Messages API 的 SSE 事件序列: +``` +event: message_start +data: {"type": "message_start", "message": {"id": "msg_01...", ...}} + +event: content_block_start +data: {"type": "content_block_start", "index": 0, "content_block": {"type": "text", "text": ""}} + +event: content_block_delta +data: {"type": "content_block_delta", "index": 0, "delta": {"type": "text_delta", "text": "Hello"}} + +event: content_block_stop +data: {"type": "content_block_stop", "index": 0} + +event: message_delta +data: {"type": "message_delta", "delta": {"stop_reason": "end_turn", "stop_sequence": null}, "usage": {"output_tokens": 50}} + +event: message_stop +data: {"type": "message_stop"} +``` + +**映射到 StreamEvent**: + +| Anthropic Event | StreamEvent | +|----------------|-------------| +| `message_start` | `MessageStart { id, model }` | +| `content_block_start { index, text }` | `ContentBlockStart { index, block_type: Text }` | +| `content_block_start { index, tool_use }` | `ContentBlockStart { index, block_type: ToolUse { id, name } }` | +| `content_block_start { index, thinking }` | `ContentBlockStart { index, block_type: Thinking }` | +| `content_block_delta { text_delta }` | `TextDelta { text }` | +| `content_block_delta { thinking_delta }` | `ThinkingDelta { text }` | +| `content_block_delta { input_json_delta }` | `ToolCallArgumentsDelta { index, arguments }` | +| `content_block_stop` | `ContentBlockEnd { index }` 或 `ToolCallEnd { index }` | +| `message_delta { delta, usage }` | `CostUpdate { usage }` + `set_thinking_signature()` (如 `delta.thinking.signature` 存在) | +| `message_stop` | 调用 `partial.finalize()` → `MessageComplete { full_response }` | +| `ping` | 忽略 | + +**注意**:与 OpenAI 不同,Anthropic 的流式有明确的 block 边界事件,映射到 `ContentBlockStart/End` 更自然。 + +### 3e:thinking_signature 处理 + +在 `message_delta` 事件中,如果 `delta.thinking?.signature` 存在,直接调用 `partial.set_thinking_signature(sig)` 将其写入内部状态。`MessageComplete` 事件不再携带 signature 字段。 + +flow 伪代码: +```rust +let mut partial = PartialMessageResponse::new(); +partial.id = Some(msg_start.message.id); +partial.model = Some(msg_start.message.model); + +while let Some(event) = stream.next().await { + match event { + AnthropicEvent::ContentBlockStart { index, content_block } => { + let block_type = match content_block.type.as_str() { ... }; + StreamEvent::ContentBlockStart { index, block_type }.apply_to(&mut partial); + } + AnthropicEvent::ContentBlockDelta { index, delta } => { + // 根据 delta.type 映射为对应 StreamEvent + let ir_event = map_delta(index, delta); + ir_event.apply_to(&mut partial); + } + AnthropicEvent::ContentBlockStop { index } => { + StreamEvent::ContentBlockEnd { index }.apply_to(&mut partial); + } + AnthropicEvent::MessageDelta { delta, usage } => { + // thinking signature: Provider 直接写入内部状态 + if let Some(sig) = delta.thinking.and_then(|t| t.signature) { + partial.set_thinking_signature(sig); + } + // CostUpdate + let partial_usage = map_usage(usage); + StreamEvent::CostUpdate { usage: partial_usage }.apply_to(&mut partial); + } + AnthropicEvent::MessageStop => { + let full = partial.finalize()?; + yield StreamEvent::MessageComplete { full_response: full }; + } + _ => {} + } +} +``` + +### 3f:capabilities() 实现 + +```rust +fn capabilities(&self) -> ProviderCapabilities { + ProviderCapabilities { + provider_name: "anthropic", + supported_models: Some(vec![ + "claude-sonnet-4-20250514".into(), + "claude-3-5-sonnet-20241022".into(), + ]), + features: ProviderFeatures { + streaming: true, + thinking: true, + vision: true, + tool_use: true, + parallel_tool_calls: true, + system_prompt_in_messages: false, // Anthropic 使用 system 参数而非 messages 中的 system role + max_context_window: 200_000, + ..Default::default() + }, + } +} +``` + +--- + +## 任务 4:实现 DeepSeekProvider / QwenProvider + +### 4a:推荐方案 — newtype 包装 GenericOpenaiProvider + +**类型别名方案(不推荐)**: +```rust +pub type DeepSeekProvider = GenericOpenaiProvider; // ❌ 零类型安全 +``` +类型别名使 `DeepSeekProvider` 和 `GenericOpenaiProvider` 完全等价——编译期无法区分, +也无法为特定 Provider 单独实现 trait(如自定义错误映射)。 + +**推荐方案 — newtype 包装**: + +将 DeepSeek 和 Qwen 合并在同一个文件 `src/llm/provider/openai_compat.rs` 中, +使用 newtype 结构体包裹 `GenericOpenaiProvider`: + +```rust +// src/llm/provider/openai_compat.rs +use crate::llm::provider::openai::GenericOpenaiProvider; + +/// DeepSeek Provider — OpenAI-compatible 协议,仅配置不同。 +pub struct DeepSeekProvider(pub GenericOpenaiProvider); + +impl DeepSeekProvider { + pub fn new(base_url: String, api_key: String, model: String) -> Self { + let url = if base_url.is_empty() { + "https://api.deepseek.com".to_string() + } else { + base_url + }; + Self(GenericOpenaiProvider::new_with_name( + url, api_key, model, "deepseek", + )) + } +} + +/// 通义千问 Provider — OpenAI-compatible 协议。 +pub struct QwenProvider(pub GenericOpenaiProvider); + +impl QwenProvider { + pub fn new(base_url: String, api_key: String, model: String) -> Self { + let url = if base_url.is_empty() { + "https://dashscope.aliyuncs.com/compatible-mode/v1".to_string() + } else { + base_url + }; + Self(GenericOpenaiProvider::new_with_name( + url, api_key, model, "qwen", + )) + } +} +``` + +**LlmProvider trait 代理**:每个 newtype 需要约 20 行 trait 代理代码。 +使用 macro 减少重复: + +```rust +macro_rules! delegate_openai_compat { + ($name:ident) => { + #[async_trait] + impl LlmProvider for $name { + async fn chat(&self, req: MessageRequest) -> Result { + self.0.chat(req).await + } + async fn chat_stream( + &self, req: MessageRequest, + ) -> Result> + Send>>, LlmError> { + self.0.chat_stream(req).await + } + fn capabilities(&self) -> ProviderCapabilities { + self.0.capabilities() + } + } + }; +} + +delegate_openai_compat!(DeepSeekProvider); +delegate_openai_compat!(QwenProvider); +``` + +### 4b:差异化处理 + +`GenericOpenaiProvider` 需要暴露配置参数支持 Provider 间的差异: + +```rust +pub struct GenericOpenaiProvider { + http_client: Client, + base_url: String, + api_key: String, + model: String, + provider_name: &'static str, + /// 额外请求头,Qwen 需要 X-DashScope-SSE: enable + extra_headers: Vec<(String, String)>, +} +``` + +| 差异点 | DeepSeek | Qwen | +|--------|----------|------| +| base_url | `api.deepseek.com` | `dashscope.aliyuncs.com/compatible-mode/v1` | +| max_tokens 字段名 | `max_tokens`(兼容) | `max_tokens`(兼容) | +| 错误格式 | 标准 OpenAI 风格 | 非标准 error body(需 fallback) | +| 模型名 | `deepseek-chat` | `qwen-plus` / `qwen-max` | +| 额外请求头 | 无特殊 | `X-DashScope-SSE: enable`(由 `extra_headers` 传入) | +| 默认模型 | `deepseek-chat` | `qwen-plus` | + +**Qwen 错误格式 fallback**: + +Qwen 的错误 body 可能是非标准 JSON 格式甚至纯文本。处理策略: + +```rust +fn handle_error(&self, status: u16, body: &str) -> LlmError { + // 先尝试按 OpenAI 标准格式解析 + if let Ok(openai_err) = serde_json::from_str::(body) { + return Self::map_openai_error(status, openai_err); + } + // fallback: 将原始文本包装为 Request 错误 + LlmError::Request { + status, + body: body.to_string(), + } +} +``` + +**实施注意**:`GenericOpenaiProvider` 的所有配置参数由 `new_with_name` 构造器统一传入, +DeepSeek/Qwen 的 newtype 构造器负责填充差异化的参数值。 + +--- + +## 任务 5:适配 ProviderRegistry + +`src/llm/provider/registry.rs` 的改动较小: +- `register_with_config()` 不需要改——它调用 `create_provider()`,而 `create_provider` 已适配新 enum +- `get()` 返回 `Option<&dyn LlmProvider>`——trait 签名变更后自动生效 +- 不需要额外修改 + +--- + +## 任务 6:添加 wiremock 集成测试 + +### 6a:Cargo.toml + +```toml +[dev-dependencies] +wiremock = "0.6" +``` + +### 6b:OpenaiChatProvider 测试 + +```rust +#[cfg(test)] +mod tests { + use super::*; + use wiremock::{Mock, MockServer, ResponseTemplate}; + use wiremock::matchers::{method, path}; + + #[tokio::test] + async fn test_openai_chat_basic() { + let mock_server = MockServer::start().await; + + // 注册 mock 响应 + Mock::given(method("POST")) + .and(path("/chat/completions")) + .respond_with(ResponseTemplate::new(200).set_body_json(json!({ + "id": "chatcmpl-123", + "object": "chat.completion", + "created": 1718000000, + "model": "gpt-4o", + "choices": [{ + "index": 0, + "message": { + "role": "assistant", + "content": "Hello!" + }, + "finish_reason": "stop" + }], + "usage": { + "prompt_tokens": 10, + "completion_tokens": 5, + "total_tokens": 15 + } + }))) + .mount(&mock_server) + .await; + + let provider = OpenaiChatProvider::new( + mock_server.uri(), "sk-test".into(), "gpt-4o".into(), + ); + + let request = MessageRequest { + model: "gpt-4o".into(), + messages: vec![Message::user_text("Hi")], + ..Default::default() + }; + + let response = provider.chat(request).await.unwrap(); + assert_eq!(response.model, "gpt-4o"); + assert_eq!(response.message.text(), "Hello!"); + } + + #[tokio::test] + async fn test_openai_chat_error() { + // HTTP 401 → LlmError::Authentication + // HTTP 429 → LlmError::RateLimit + // HTTP 500 → LlmError::Request + } + + #[tokio::test] + async fn test_openai_chat_stream() { + // SSE 流式响应测试 + } +} +``` + +### 6c:AnthropicProvider 测试 + +```rust +#[tokio::test] +async fn test_anthropic_chat_basic() { + let mock_server = MockServer::start().await; + + // Mock Anthropic Messages API 响应 + Mock::given(method("POST")) + .and(path("/v1/messages")) + .respond_with(ResponseTemplate::new(200).set_body_json(json!({ + "id": "msg_01...", + "type": "message", + "role": "assistant", + "content": [{"type": "text", "text": "Hello from Claude!"}], + "model": "claude-sonnet-4-20250514", + "stop_reason": "end_turn", + "usage": {"input_tokens": 10, "output_tokens": 5} + }))) + .mount(&mock_server) + .await; + + let provider = AnthropicProvider::new( + mock_server.uri(), "sk-ant-test".into(), "claude-sonnet-4-20250514".into(), + ); + + let request = MessageRequest { + model: "claude-sonnet-4-20250514".into(), + messages: vec![Message::user_text("Hi")], + max_tokens: Some(100), + ..Default::default() + }; + + let response = provider.chat(request).await.unwrap(); + assert_eq!(response.message.text(), "Hello from Claude!"); +} + +#[tokio::test] +async fn test_anthropic_chat_stream() { + // SSE 事件序列测试 + // message_start → content_block_start(text) → content_block_delta(text_delta) + // → content_block_stop → message_delta → message_stop +} +``` + +### 6d:Provider 双向映射测试 + +测试 `Message → OpenaiChatRequest → MessageResponse` 和 `Message → AnthropicRequest → MessageResponse` 的双向转换一致性。 + +### 6e:测试覆盖要求 + +每个 Provider 的测试至少覆盖以下场景: + +**基本路径**: +- 基本文本对话(chat + chat_stream) +- 带工具定义的对话 +- 多轮消息历史(system + user + assistant + tool_result) +- 空消息列表 → 合理错误 + +**错误路径**: +- HTTP 401 → `LlmError::Authentication` +- HTTP 429 → `LlmError::RateLimit` +- HTTP 500 → `LlmError::Request` +- HTTP 529(Anthropic overloaded)→ `LlmError::RateLimit` + +**边界场景**: +- 响应中 `choices` 数组为空 → `convert_response` 返回错误而非 panic +- 流式 `choices: []` 的中间 chunk → 忽略,不产生事件 +- SSE 中断后部分数据 → 解析不 panic,返回残数据错误 +- `data: [DONE]` 之前流中断 → 产出部分响应 +- 并行 tool_calls(多 index 同时出现在同一切片)→ 正确处理 index 分配 + +**Anthropic 特有**: +- `message_start` → `content_block_start(text)` → `content_block_delta(text_delta)` → `content_block_stop` → `message_delta` → `message_stop` 完整序列 +- `ping` 事件 → 忽略,不中断流 +- `message_start` 后直接 `message_stop`(空响应)→ 产出空 MessageResponse + +--- + +## 验证方式 + +1. **编译检查**:`cargo build` 通过 +2. **单元测试**:`cargo test` 全部通过 +3. **Provider 测试**:每个 Provider 的 mock 集成测试覆盖基本路径 + 错误路径 +4. **clippy 检查**:`cargo clippy` 无新增警告 + +## 回滚方案 + +如果某个 Provider 实现不合理,将该 Provider 回退为 `Err(LlmError::Other(...))`,不影响其他 Provider。每个 Provider 完成时建议打 tag `provider-{name}-v2`。 + +## 开放事项 + +- **`OpenaiResponseProvider` 完整实现**:当前返回 `Err(LlmError::Other(...))`。升级触发条件:Trace 中 Response API 采用率 >20%,或收到 ≥3 个企业客户需求。升级后覆盖核心对话能力(create response + 流式 + 工具调用),内置工具(web_search, file_search)通过 `MessageRequest.extra` 传递。 +- **SSE 解析通用层**:当前 `OpenaiProvider` 有内联 `SseChunkStream`,`AnthropicProvider` 也有类似 SSE 解析逻辑。是否在 Phase 1 提取通用 `SseByteStream` 层(独立于类型系统、纯字节处理),在"Phase 0 未完成时的降级路径"中可作为先行完成项。 +- **`ContentBlock::Extension` 在 OpenAI Response API 内置工具场景的使用** +- **Google Gemini Provider**:当前不在 Phase 1 范围内。建议的触发条件:3 个客户需求或 Gemini API 流量占项目总 LLM 调用的 5% 以上时进入"Next"阶段。 +- **Anthropic `anthropic-beta` header 与 Prompt Caching**:Prompt caching 等特性需通过此 header 启用。当前 Phase 1 不做,当用户明确要求缓存或流量成本显著时启动。 diff --git a/docs/10c-phase2-llm-cycle-simplify.md b/docs/10c-phase2-llm-cycle-simplify.md new file mode 100644 index 0000000..386a856 --- /dev/null +++ b/docs/10c-phase2-llm-cycle-simplify.md @@ -0,0 +1,1049 @@ +# Phase 2 实施计划:LlmCycle 简化(逻辑重构) + +> **所属方案**:[10-llm-provider-refinement.md](10-llm-provider-refinement.md) +> +> **前置条件**:Phase 0 + Phase 1 已完成——新类型已就位,所有 Provider 已适配新 trait 签名 +> +> **⚠️ 交叉阶段修复**:审查发现部分关键修复需在 Phase 0 期间处理,否则 Phase 2 无法直接实施。详见下文"审查发现与交叉阶段修复计划"。 +> +> **目标**:移除 Phase 0 遗留的临时转换层,将 `LlmCycle` 内部消息存储从 `Vec` 切换为 `Vec`,利用新类型的表达能力重写核心逻辑 + +--- + +## 目标 + +1. 将 `LlmCycle` 内部消息存储从 `Vec` 切换为 `Vec`,**移除 Phase 0 引入的 `OpenaiChatMessage ←→ Message` 转换层** +2. 流处理循环重构:利用 `MessageComplete.full_response` 直接拿到完整响应,去掉手动 delta 累积 +3. 工具循环清洗:从 `MessageResponse.message` 的 content 中直接提取 `ContentBlock::ToolUse` +4. `compact.rs` 适配新 `Message` 类型 +5. `ConversationMemory` 的消息类型同步切换 + +--- + +## 审查发现与交叉阶段修复计划 + +> 本方案于 2026-06-30 经架构审查发现 4 个必须修复的问题和多个需在 Phase 0 提前处理的依赖项。以下修复计划按"所属阶段"分类,标注责任文件和建议修复方式。 + +### 🔴 交叉阶段修复(Phase 0 期间完成,否则 Phase 2 受阻) + +#### FIX-A:`ToolInvocation` 增加 `tool_call_id` 字段 + +| 属性 | 值 | +|------|-----| +| **严重程度** | 🔴 阻碍性 | +| **影响** | 工具结果回传时 tool_call_id 丢失,Anthropic Provider 上必然失败 | +| **文件** | `src/tools/registry.rs`(`ToolInvocation` 结构体)、`src/llm/cycle.rs`(调用点) | + +**问题跟踪**: +- `ToolInvocation`(`registry.rs:17-35`)只有 `tool_name`、`input`、`output` 三个字段,**没有 `tool_call_id`** +- `invoke_all()` 签名(`registry.rs:149-151`)接收 `Vec<(String, Value)>`(name, args),**未传递 id** +- `invoke()` 签名(`registry.rs:131`)是 `(&self, name: &str, args: Value)`,**未传递 id** +- 当前代码用 `result.tool_name` 冒充 `tool_call_id`(`cycle.rs:617`)——对 OpenAI 碰巧可用,Anthropic 必然失败 +- `has_tool_calls_in_message`(`cycle.rs:645`)基于旧 `OpenaiChatMessage::Assistant { tool_calls }` 判断,Phase 2 需同步改为检查 `ContentBlock::ToolUse` + +**修复方案**: + +```rust +// ToolInvocation 增加 tool_call_id(必选,每次调用必有真实值) +pub struct ToolInvocation { + pub tool_call_id: String, // 新增 + pub tool_name: String, + pub input: Value, + pub output: Result, +} + +// ToolContext 不增加 tool_call_id: +// 当前 ToolContext::new(name, "") 中 trace_id 已被设为空串, +// tool_call_id 通过 invoke_all→invoke→ToolInvocation 传递即可, +// 工具实现方不需要感知 tool_call_id。 +``` + +> **关于 `ToolContext` 的决策**:不向 `ToolContext` 增加 `tool_call_id`。理由是 `ToolContext` 目前通过 `invoke()` 内部构造(`registry.rs:140`:`ToolContext::new(name, "")`),传递的是 `tool_name` 而非 `tool_call_id`。若增加 `tool_call_id` 字段会要求所有 `invoke` 调用链都感知 id——这增加了复杂度,而工具实现方实际并不需要 tool_call_id(结果是按调用顺序返回的,工具通过在消息历史中搜索 `ToolUse` 来关联上下文)。**`tool_call_id` 仅在 `ToolInvocation` 中保存,供 `submit_with_tools` 回传时使用。** + +**数据流**:`extract_tool_calls_from_message()` 返回 `(id, name, args)` → `invoke_all` 签名改为 `Vec<(String, String, Value)>`(id, name, args)→ `ToolInvocation.tool_call_id` → `submit_with_tools()` 回传工具结果时使用 `result.tool_call_id`。`invoke()` 内部在构造 `ToolInvocation` 时即可获得 tool_call_id。 + +**建议归属**:Phase 2.2(随 Task 4 统一改 `registry.rs` + `cycle.rs`,避免切割到 Phase 0 导致 scope 漂移) + +**影响范围**: +- `ToolInvocation::new()` 签名从 3 参数变为 4 参数(`registry.rs:28`) +- `invoke()` 签名从 `(&self, name, args)` 改为 `(&self, tool_call_id, name, args)` 或内部通过包装传递(`registry.rs:131`) +- `invoke_all()` 签名从 `Vec<(String, Value)>` 改为 `Vec<(String, String, Value)>`(`registry.rs:149-151`) +- `cycle.rs:584-594` 中 `|(_id, name, args)|` → `|(id, name, args)|` 不再丢弃 id +- 测试代码中的 `ToolInvocation` 和 `invoke_all` 调用点同步更新 + +--- + +#### FIX-B:`Message` + `ContentBlock` 的 serde 派生在 Phase 0 添加 + +| 属性 | 值 | +|------|-----| +| **严重程度** | 🔴 阻碍性 | +| **影响** | Phase 0→Phase 2 过渡期 ConversationMemory 无法使用新类型 | + +**问题**:Phase 2 的 Task 7(ConversationMemory 适配)依赖 `Message` 实现 `Serialize`/`Deserialize`。如果 Phase 0 定义 `Message` 时未添加 serde 派生,Phase 2 需要回头修改 Phase 0 创建的文件,造成"谁创建谁维护"责任混乱。 + +**修复方案**:Phase 0 定义新类型时直接添加 serde 派生: + +```rust +#[derive(Debug, Clone, Serialize, Deserialize)] +#[serde(tag = "type", rename_all = "snake_case")] +pub enum Message { + System { content: Vec }, + User { content: Vec }, + UserImage { data: String, mime_type: String, detail: ImageDetail }, + Assistant { content: Vec }, + ToolResult { tool_call_id: String, content: Vec, is_error: bool }, +} +``` + +`ContentBlock` 同样添加 serde 派生(tag 策略在 Phase 0 设计时决定,建议 `#[serde(rename_all = "snake_case")]` + `tag = "type"` 或 untagged 视用法而定)。 + +此序列化格式**仅用于内部存储**(ConversationMemory 持久化),不与任何 Provider 协议交互,故无需考虑外部兼容性。`UserImage` 的 JSON 形如 `{"type": "user_image", "data": "...", "mime_type": "image/png", "detail": "auto"}`。 + +**同步添加构造方法**(`impl Message` 块中的便捷构造函数,FIX-D 的 deprecated shim 和 Phase 2 各 Task 都依赖它们): + +```rust +impl Message { + pub fn system_text(text: impl Into) -> Self { + Message::System { content: vec![ContentBlock::Text { text: text.into() }] } + } + pub fn user_text(text: impl Into) -> Self { + Message::User { content: vec![ContentBlock::Text { text: text.into() }] } + } + pub fn assistant(text: impl Into) -> Self { + Message::Assistant { content: vec![ContentBlock::Text { text: text.into() }] } + } + pub fn tool_result(tool_call_id: impl Into, text: impl Into, is_error: bool) -> Self { + Message::ToolResult { + tool_call_id: tool_call_id.into(), + content: vec![ContentBlock::Text { text: text.into() }], + is_error, + } + } +} +``` + +**建议归属**:Phase 0,作为类型定义的一部分(零成本加法) + +--- + +#### FIX-C:`ContentBlock` 变体列表在 Phase 0 设计时明确 + +| 属性 | 值 | +|------|-----| +| **严重程度** | 🔴 阻碍性 | +| **影响** | 方案中 `estimate_block_tokens` 的 match 分支引用了未定义的 `ContentBlock::ToolResult`,导致 Phase 2 实施时类型不匹配 | + +**问题**:方案 §6b 的 `estimate_block_tokens` 中出现 `ContentBlock::ToolResult { content, .. }` 分支(第 374 行),但父方案 `10-llm-provider-refinement.md` §2.1 定义 `ContentBlock` 时并未包含 `ToolResult` 变体。`ToolResult` 是 `Message` 的变体(消息级角色),不是 `ContentBlock` 的变体(内容块级角色)。 + +**分析**:在 Anthropic 协议中,`Assistant` 消息的 content 可以直接包含 `tool_result` block——这对应 `ContentBlock::ToolResult` 的应用场景(嵌套在其他消息的 content 中);而 `Message::ToolResult` 对应 OpenAI 的 `role: "tool"` 消息(消息列表的顶层变体)。二者语义层级不同,**可以合理共存**。 + +**修复方案**:Phase 0 设计 `ContentBlock` 时明确是否包含 `ToolResult` 变体。推荐包含,因为: + +1. Anthropic 协议支持 Assistant 消息内嵌 tool_result block +2. 简化 Phase 2 的 `estimate_block_tokens` 实现(不需要在 match 中分流) +3. 让 `compact` 估算函数能递归处理嵌套的 tool result content + +Phase 2 的 `estimate_block_tokens` 改为穷尽 match(不再使用 `_ => 50` 兜底),编译器强制检查所有变体覆盖。 + +**建议归属**:Phase 0 设计决策,Phase 2 代码编写时确认 + +--- + +#### FIX-D:`system_prompt` 字段彻底移除 + +| 属性 | 值 | +|------|-----| +| **严重程度** | 🔴 阻碍性 | +| **影响** | 与父方案 §2.4 设计意图矛盾,"两条路径"(字段 + messages 中的 System 消息)可能导致重复/覆盖 | + +**问题**:方案 §2a 保留 `system_prompt: Option` 字段,在 `build_request()` 中转为 `Message::System` 插入 messages 开头。但父方案 §2.4 明确意图:"system prompt 通过 `Message::System` 在 messages 中表达,Provider 映射层自行处理差异"。保留字段导致两条路径: + +- 路径 1:调用方通过 `with_system_prompt()` 设置纯文本 +- 路径 2:调用方直接通过 `with_messages()` 传入 `Message::System` + +`build_request()` 中通过 `!messages.iter().any(|m| matches!(m, Message::System { .. }))` 尝试检测"是否已有系统消息",但本质上无法区分"用户主动传入的 System"和"之前阶段插入的 System"。 + +**修复方案**: +1. **彻底移除** `system_prompt: Option` 字段 +2. **移除** `with_system_prompt()` 方法 +3. 调用方改为使用 `Message::system_text()` + `with_messages()` 或 builder 模式 +4. `build_request()` 直接使用 `self.messages`,不做任何 System 消息插入 + +```rust +pub struct LlmCycle { + provider: Arc, + config: CycleConfig, + usage: CostTracker, + messages: Vec, // 调用方自行管理所有消息类型 + // system_prompt 已移除 + hook_executor: Option>, + compact_config: Option, + compact_state: CompactState, +} +``` + +如需向后兼容(有已有调用方使用 `with_system_prompt`),可在 `LlmCycle::new()` 或 `submit()` 前设置一个短暂的弃用期: + +```rust +#[deprecated(since = "0.x.0", note = "请改用 Message::system_text() + with_messages()")] +pub fn with_system_prompt(mut self, prompt: String) -> Self { + self.messages.insert(0, Message::system_text(prompt)); + self +} +``` + +**建议归属**:Phase 2.0(与 Task 1 消息切换一步到位,杜绝过渡期。不在 Phase 2.1 留存已废弃字段) + +**调用方迁移**(需在 Phase 2.0 实施前确认): + +| 调用位置 | 代码 | 迁移方式 | +|---------|------|---------| +| `src/agent/session.rs:141-143` | `if let Some(prompt) = self.agent.system_prompt() { cycle = cycle.with_system_prompt(prompt.to_string()); }` | 改为 `cycle = cycle.with_messages(vec![Message::system_text(prompt)])`,或在 `submit_turn` 内部直接 push | + +**依赖**:`Message::system_text()` 构造方法必须在 Phase 0 定义 `Message` 类型时同步添加(见 FIX-B)。`#[deprecated]` shim 中调用的 `Message::system_text(prompt)` 依赖于它。 + +--- + +### 🟡 Phase 2 内部修复 + +#### FIX-E:`submit_stream` 返回类型修正 + 闭包借用问题解决 + +| 属性 | 值 | +|------|-----| +| **严重程度** | 🟡 必须修复(编译错误 + API 语义断裂) | +| **影响** | 方案代码存在编译错误(yield Ok(event) vs 签名不匹配);`&mut self` 借用冲突;post_request hook 消息不完整 | + +**三个问题**: + +1. **返回类型不匹配**(§3a 第 159 行):`yield Ok(event.clone())` 产生 `Result`,但函数签名声明 `Item = StreamEvent`。 + +2. **`&mut self` 借用冲突**(§3a 第 163 行):`self.messages.push(...)` 在 `stream!` 闭包中需要 `&mut self`,但闭包只能捕获共享引用。 + +3. **post_request hook 消息不完整**(§3a 第 149 行):`let post_request = self.build_request(&tools)` 在流启动前执行,此时 `self.messages` 只有 user 消息,不包括即将到来的 Assistant 响应。 + +**⚠️ 审查结论:mpsc 方案不可行,采用简化方案作为主方案** + +**为什么 mpsc 方案不成立**(两次审查确认): + +| 问题 | 描述 | 根本原因 | +|------|------|---------| +| `rx` 竞争 | mpsc `rx` 只能有一个消费者,但代码中两个 `async move` 块都试图消费 | 编译错误 | +| 阻塞破坏流契约 | `response_rx.await` 在返回流前阻塞等待整个流结束 | 不是流式 API | +| `StreamEvent::StreamEnd` 不存在 | 引用了未定义的变体 | 新旧 StreamEvent 均无此变体 | +| `tokio::spawn` task leak | JoinHandle 被丢弃,调用方 drop Stream 后 task 可能永久挂起 | 资源泄漏 | + +**核心矛盾**:`stream!` 宏产生的异步生成器无法持有 `&mut self`(Pin 限制),且流是"延迟求值"的——`submit_stream` 返回流对象时流尚未执行。这意味着**没有任何通道方案能让 `submit_stream` 在返回流之前就拿到完整响应**。如果想在流结束后更新 `self.messages`,必须让调用方在流结束后显式调用一个方法。 + +**采用简化方案作为主方案**: + +```rust +pub async fn submit_stream( + &mut self, + prompt: String, + tools: Vec, +) -> Result + Send>>, LlmError> { + self.messages.push(Message::user_text(prompt)); + self.maybe_compact(); + + let request = self.build_request(&tools); + + // PreRequest hook + if let Some(ref executor) = self.hook_executor { + let ctx = HookContext::new(HookEvent::PreRequest).with_request(&request); + let results = executor.execute(HookEvent::PreRequest, &ctx).await; + if results.iter().any(|r| r.should_block) { + return Err(LlmError::Other("Blocked by pre-request hook".to_string())); + } + } + + let event_stream = self.provider.chat_stream(request).await?; + let hook_executor = self.hook_executor.clone(); + + Ok(Box::pin(stream! { + use futures_util::StreamExt; + let mut event_stream = event_stream; + + while let Some(result) = event_stream.next().await { + match result { + Ok(event) => { + yield event; + } + Err(e) => { + if let Some(ref executor) = hook_executor { + let ctx = HookContext::new(HookEvent::OnError).with_error(&e); + executor.execute(HookEvent::OnError, &ctx).await; + } + yield StreamEvent::Error { message: e.to_string() }; + break; + } + } + } + + // PostRequest hook:注意此处的消息列表不包含本次 Assistant 响应 + // (因为在流启动前 request 就已构建完成) + if let Some(ref executor) = hook_executor { + let ctx = HookContext::new(HookEvent::PostRequest); + executor.execute(HookEvent::PostRequest, &ctx).await; + } + })) +} +``` + +**变更**: +1. `yield event` 直接产生 `StreamEvent`(与签名 `Item = StreamEvent` 一致)——移除 `Ok()` 包装 +2. 不再在闭包中操作 `self.messages`——调用方在收到 `MessageComplete` 事件后自行调用 `cycle.push_message(full_response.message)` +3. post_request hook 明确文档化限制 + +**调用方使用模式**: + +```rust +let mut stream = cycle.submit_stream(prompt, tools).await?; +let mut final_response: Option = None; + +while let Some(event) = stream.next().await { + match &event { + StreamEvent::MessageComplete { full_response } => { + final_response = Some(full_response.clone()); + } + // 处理其他事件... + _ => {} + } +} + +// 流结束后更新消息历史 +if let Some(resp) = final_response { + cycle.push_message(resp.message.clone()); + // usage update via cycle.usage().add(&resp.usage); +} +``` + +**建议归属**:Phase 2.2,Task 3 实施时 + +--- + +#### FIX-F:`microcompact` 跳过 `is_error: true` 的工具结果 + +| 属性 | 值 | +|------|-----| +| **严重程度** | 🟡 建议修复 | +| **影响** | 错误结果包含对 LLM 理解失败原因至关重要的信息,压缩后 LLM 无法理解 | + +**问题**:§6c 的 `microcompact()` 对所有 `Message::ToolResult` 无条件压缩,包括 `is_error: true` 的结果。 + +**修复方案**:只压缩 `is_error: false` 的结果: + +```rust +// microcompact 压缩阶段 +for msg in &messages[..prune_start] { + if matches!(msg, Message::ToolResult { is_error: false, .. }) { + freed_tokens += estimate_single_message_tokens(msg); + } +} + +// 替换内容 +for msg in &mut messages[..prune_start] { + if let Message::ToolResult { content, is_error: false } = msg { + *content = vec![ContentBlock::Text { text: "[pruned]".to_string() }]; + } +} +``` + +**建议归属**:Phase 2,Task 6 实施时一并处理 + +--- + +#### FIX-G:`submit_with_tools()` 中 `extract_tool_calls_from_response` 的 `input` 类型修正 + +| 属性 | 值 | +|------|-----| +| **严重程度** | 🟡 编译类型安全(`serde_json::to_string` 在 `&Value` 上调用) | +| **影响** | 编译通过,但 serde_json::to_string(&input) 在 `&Value` 上调用可能产生意外格式 | + +**问题**:§4a 第 231-233 行 `extract_tool_calls_from_response` 中的 `input` 是 `&Value`(来自 `ContentBlock::ToolUse { input, .. }`),调用 `serde_json::to_string(input)` 可能序列化为多余的引号包裹格式。 + +**修复方案**:使用 `serde_json::to_string(input).unwrap_or_default()` 没问题(`Value` 实现了 `Serialize`),但建议改为 `input.to_string()`(Value 的 Display 实现生成紧凑 JSON): + +```rust +let args = input.to_string(); // 或 serde_json::to_string(input).unwrap_or_default() +``` + +**建议归属**:Phase 2,Task 4 实施时处理 + +--- + +#### FIX-H:测试基础设施迁移清单 + +| 属性 | 值 | +|------|-----| +| **严重程度** | 🟡 必须覆盖(否则测试无法编译) | +| **影响** | 测试代码使用旧类型,Phase 2 实施后编译失败 | + +**需迁移的测试辅助代码**(按文件列出): + +| 文件 | 行号 | 实体 | 当前类型 | 目标类型 | +|------|------|------|---------|---------| +| `src/llm/cycle.rs` | 700-725 | `MockProvider` | `ChatRequest`/`ChatResponse` | `MessageRequest`/`MessageResponse` | +| `src/llm/cycle.rs` | 731-737 | `assistant_text_response()` | 返回 `ChatResponse` | 返回 `MessageResponse` | +| `src/llm/cycle.rs` | 739-765 | `assistant_tool_call_response()` | `OpenaiToolCall` 构造 | `ContentBlock::ToolUse` 构造 | +| `src/agent/session.rs` | 202-224 | `MockProvider` | `ChatRequest`/`ChatResponse` | `MessageRequest`/`MessageResponse` | +| `src/agent/session.rs` | 241-247 | `assistant_text()` | 返回 `ChatResponse` | 返回 `MessageResponse` | +| `src/agent/session.rs` | 141-143 | `with_system_prompt()` 调用点 | `OpenaiChatMessage` | `Message::system_text()` | +| `src/agent/builder.rs` | 126-132 | `StubProvider` | `ChatRequest`/`ChatResponse` | `MessageRequest`/`MessageResponse` | +| `src/tools/registry.rs` | 28 | `ToolInvocation::new()` | 3 参数 | 4 参数(含 `tool_call_id`) | + +**建议归属**:Phase 2 尾部,作为一个单独的"测试适配 PR"集中修改(不在各 Task 中随改),避免业务代码和测试代码混杂。Phase 2 业务代码合入后,紧接着提交测试适配 PR。 + +--- + +### 🔄 实施顺序调整 + +原方案的线性顺序(1→2→3→4→5→6→7→8→9)考虑了依赖关系,但存在优化空间。结合交叉修复后,推荐以下实施顺序: + +``` +[Phase 0] 类型层(含交叉修复) + FIX-B Message + ContentBlock serde 派生 + 构造方法 ← 类型定义时添加 + FIX-C ContentBlock 变体列表明确 ← 设计决策 + +[Phase 2.0] 基础设施切换 + Task 1 LlmCycle 消息切换 (Vec → Vec) + 含 FIX-D system_prompt 一步到位移除(与消息类型切换同时做) + Task 7 ConversationMemory 切换 ← 可并行,依赖 FIX-B + +[Phase 2.1] 核心逻辑简化 + Task 5 submit/submit_messages/submit_request ← 依赖 Phase 2.0,最简单先做 + Task 2 build_request 简化 ← 依赖 Phase 2.0(system_prompt 已在 Phase 2.0 移除) + +[Phase 2.2] 流处理 + 工具循环 + Task 3 submit_stream 简化 ← 简化方案(非 mpsc),含 FIX-E + Task 4 工具循环清洗 + FIX-A(tool_call_id) ← 统一改 registry + cycle + +[Phase 2.3] 附属适配 + 清理 + Task 6 compact.rs 适配 ← 可并行,含 FIX-F + Task 9 prompt/composer.rs 适配 + Task 8 清理临时转换函数(含 has_tool_calls_in_message) + +[Phase 2.4] 收尾 + FIX-H 测试基础设施迁移 ← 单独 PR 集中处理 + (可选) FIX-G 代码风格优化 +``` + +--- + +| 操作 | 文件 | 说明 | +|------|------|------| +| 修改 | `src/llm/cycle.rs` | 主要修改:消息类型切换 + 逻辑简化 | +| 修改 | `src/llm/compact.rs` | 适配 `Message` 类型(`microcompact` / `should_compact` / 估算函数) | +| 修改 | `src/memory/conversation.rs` | 消息类型从 `OpenaiChatMessage` 切换为 `Message` | +| 修改 | `src/prompt/composer.rs` | `validate_messages()` 适配 `Message` 类型 | +| 清理 | 各处 | 删除 Phase 0 引入的临时转换函数 | + +**不涉及**: +- `src/llm/cycle/usage.rs` — `Usage` 类型不变 +- `src/llm/cycle/retry.rs` — `RetryConfig` 和 `should_retry` 不引用消息类型 +- `src/llm/error.rs` — `LlmError` 不引用消息类型 + +--- + +## Phase 2 对外影响(breaking changes 汇总) + +以下变更影响所有使用 `LlmCycle` 的上游代码。Phase 2 实施前需确认调用方已适配。 + +| 变更 | 旧方式 | 新方式 | 影响文件 | +|------|--------|--------|---------| +| `messages()` 返回类型 | `&[OpenaiChatMessage]` | `&[Message]` | 所有读取消息历史的上游代码 | +| `with_messages()` / `extend_messages()` 参数 | `Vec` | `Vec` | 使用 builder 模式的代码 | +| `with_system_prompt()` | 设置纯文本,`build_request()` 隐式插入 | **已废弃**(`#[deprecated]`),调用方使用 `Message::system_text()` + `with_messages()` | `src/agent/session.rs:141-143` | +| `submit()` 返回类型 | `Result` | `Result` | 所有调用 `submit()` 的代码 | +| `submit_with_tools()` 返回类型 | `Result` | `Result` | `AgentSession` | +| `submit_messages()` 签名 | `Vec` 参数 | `Vec` 参数 | 直接调用方 | +| `submit_stream()` 调用方式 | 流内回传 `StreamEvent`(旧),`submit_stream` 负责 `self.messages.push` | `StreamEvent`(新)透传,调用方需在 `MessageComplete` 后手动调用 `cycle.push_message()` | 调用流式 API 的代码 | +| `compact` 模块接收 | `&[OpenaiChatMessage]` 参数 | `&[Message]` 参数 | 仅在 `LlmCycle` 和 `ConversationMemory` 内部调用,对外无影响 | + +**迁移建议**:所有上游变更集中在 `src/agent/session.rs` 和测试代码中。建议 Phase 2.0 实施前先在 `session.rs` 中完成消息类型对齐,然后随 Phase 2 各阶段逐步适配。 + +--- + +## 任务 1:LlmCycle 内部消息切换 + +### 1a:结构体字段变更 + +```rust +pub struct LlmCycle { + provider: Arc, + config: CycleConfig, + usage: CostTracker, + // 从 Vec 改为 Vec + messages: Vec, + // system_prompt 字段已移除(审查 FIX-D): + // 调用方通过 Message::system_text() + with_messages() 管理系统提示 + // 向后兼容:with_system_prompt() 标记为 #[deprecated],内部转为 Message::System 后插入 + hook_executor: Option>, + compact_config: Option, + compact_state: CompactState, +} +``` + +### 1b:`messages()` 方法签名变更 + +```rust +// 当前 +pub fn messages(&self) -> &[OpenaiChatMessage] { &self.messages } +// 目标 +pub fn messages(&self) -> &[Message] { &self.messages } +``` + +### 1c:`with_messages()` 方法签名变更 + +```rust +// 当前 +pub fn with_messages(mut self, messages: Vec) -> Self +// 目标 +pub fn with_messages(mut self, messages: Vec) -> Self +``` + +### 1d:`extend_messages()` 方法签名变更 + +```rust +// 当前 +pub fn extend_messages(&mut self, messages: Vec) +// 目标 +pub fn extend_messages(&mut self, messages: Vec) +``` + +### 1e:新增 `push_message()` 公共方法 + +供 `submit_stream()` 消费方在收到 `MessageComplete` 事件后手动追加消息: + +```rust +/// 追加单条消息到历史尾部(供 submit_stream 消费方使用)。 +pub fn push_message(&mut self, msg: Message) { + self.messages.push(msg); +} +``` + +--- + +## 任务 2:简化 `build_request()` + +### 2a:移除转换步骤 + 移除 system_prompt 逻辑(审查 FIX-D) + +Phase 0 中的 `build_request()` 需要调用 `chat_message_to_message` 将 `self.messages`(`Vec`)转为 `Vec`。Phase 2 中 `self.messages` 本身就是 `Vec`,所以这个转换步骤被移除。 + +同时,`system_prompt` 字段已移除(审查 FIX-D),调用方通过 `Message::system_text()` + `with_messages()` 自主管理系统提示消息。 + +```rust +fn build_request(&self, tools: &[ToolDefinition]) -> MessageRequest { + // 直接使用 self.messages,不做任何 System 消息隐式插入 + let messages = self.messages.clone(); + + MessageRequest { + model: self.config.model.clone(), + messages, + tools: tools.to_vec(), + tool_choice: ToolChoice::Auto, + max_tokens: self.config.max_tokens, + temperature: self.config.temperature, + top_p: None, + stop_sequences: vec![], + stream: false, + thinking: None, + extra: HashMap::new(), + } +} +``` + +> **注意**:`tools` 的类型从 `Option>`(OpenAI 特定包装)变为 `Vec`(`MessageRequest` 中使用的新类型)。Provider 映射层负责将 `ToolDefinition` 转换为协议特定的工具定义格式。`build_request()` 不再需要 `OpenaiTool::Function` 包装。 + +### 2b:`with_system_prompt()` 弃用过渡 + +```rust +#[deprecated(since = "0.x.0", note = "请改用 Message::system_text() + with_messages()")] +pub fn with_system_prompt(mut self, prompt: String) -> Self { + self.messages.insert(0, Message::system_text(prompt)); + self +} +``` + +此方法保留一个过渡周期(标注 `#[deprecated]`),编译器发出弃用警告但不阻止编译。后续版本中移除。 + +### 2c:`with_messages()` 中忽略 system_prompt 冲突 + +`with_messages()` 不再做去重检查或 System 消息过滤——调用方传入什么就是什么。`build_request()` 直接使用 `self.messages`,不插入任何额外消息。 + +--- + +## 任务 3:简化 `submit_stream()` + +### 3a:流处理循环变更(审查 FIX-E — 简化方案) + +Phase 0 中,`submit_stream()` 从 `provider.chat_stream()` 直接获得 `Result` 流。但为了在流结束后仍能向 `self.messages`(`Vec`)追加响应,需要通过 `message_to_chat_message` 转换。 + +Phase 2 中,`self.messages` 已经是 `Vec`,所以不需要转换了。 + +> ⚠️ **审查结论**:mpsc 通道方案因 Rust 的 Pin/借用限制不可行(`&mut self` 无法进入闭包,流是延迟求值的),采用**简化方案**——`submit_stream` 内部不管理消息历史,由调用方在流结束后自行调用 `push_message()`。 +> +> 详情见"审查发现与交叉阶段修复计划"中 FIX-E 的分析。 + +```rust +pub async fn submit_stream( + &mut self, + prompt: String, + tools: Vec, +) -> Result + Send>>, LlmError> { + self.messages.push(Message::user_text(prompt)); + self.maybe_compact(); + + let request = self.build_request(&tools); + + // PreRequest hook + if let Some(ref executor) = self.hook_executor { + let ctx = HookContext::new(HookEvent::PreRequest).with_request(&request); + let results = executor.execute(HookEvent::PreRequest, &ctx).await; + if results.iter().any(|r| r.should_block) { + return Err(LlmError::Other("Blocked by pre-request hook".to_string())); + } + } + + let event_stream = self.provider.chat_stream(request).await?; + let hook_executor = self.hook_executor.clone(); + + Ok(Box::pin(stream! { + use futures_util::StreamExt; + let mut event_stream = event_stream; + + while let Some(result) = event_stream.next().await { + match result { + Ok(event) => { + yield event; + } + Err(e) => { + if let Some(ref executor) = hook_executor { + let ctx = HookContext::new(HookEvent::OnError).with_error(&e); + executor.execute(HookEvent::OnError, &ctx).await; + } + yield StreamEvent::Error { message: e.to_string() }; + break; + } + } + } + + // PostRequest hook:注意此处的消息列表不包含本次 Assistant 响应 + if let Some(ref executor) = hook_executor { + let ctx = HookContext::new(HookEvent::PostRequest); + executor.execute(HookEvent::PostRequest, &ctx).await; + } + })) +} +``` + +**调用方使用模式**: + +```rust +let mut stream = cycle.submit_stream(prompt, tools).await?; +let mut final_response: Option = None; + +while let Some(event) = stream.next().await { + match &event { + StreamEvent::MessageComplete { full_response } => { + final_response = Some(full_response.clone()); + } + // 其他事件(TextDelta, ContentBlockStart 等)透传给 UI/消费方 + _ => {} + } +} + +// 流结束后更新消息历史 +if let Some(resp) = final_response { + cycle.push_message(resp.message.clone()); + // usage 更新:cycle.usage() 返回 &CostTracker,调用方自行 add +} +``` + +**变更**: +- `yield` 直接产生 `StreamEvent`(签名 `Item = StreamEvent`),不再包裹 `Ok()` +- `self.messages.push(response.message)` 由调用方在流结束后通过 `push_message()` 完成 +- `self.usage.add(&response.usage)` 由调用方通过 `cycle.usage().add()` 完成 +- 无需 `message_to_chat_message` 转换 + +### 3b:设计决策记录 + +**决策**:采用简化方案(已验证 mpsc 方案因 Rust 语言限制不可行)。 + +**已知限制**(调用方需知晓): +- post_request hook 收到的消息列表中**不包含本次响应**——如需完整消息列表,调用方应在流结束后手动触发 Hook + +**判断标准**(如后续需切换方案): +- 如果未来有需求要求 `submit_stream` 在流结束后自动推消息到 `self.messages`,可考虑将 `self.messages` 包裹为 `Arc>>` 并在流闭包内安全修改——但这是 API 层面的 breaking change + +--- + +## 任务 4:清洗工具循环 + +### 4a:`submit_with_tools()` 中的 tool 循环 + +当前 `submit_with_tools()` 使用 Phase 0 的临时逻辑: +- `response.message` 是 `Message` 类型(新枚举) +- 使用 `has_tool_calls_in_response` 和 `extract_tool_calls_from_response` 解析 tool calls +- 使用 `OpenaiChatMessage::tool_result()` 推送 tool 结果到 `self.messages` + +Phase 2 中: +- `self.messages.push(response.message.clone())` —— 直接推送 `Message`(不需要转换) +- tool 结果回传:`self.messages.push(Message::tool_result(id, content, is_error))` +- 从 `response.message` 提取 tool calls:遍历 `content` 中的 `ContentBlock::ToolUse` + +```rust +fn has_tool_calls_in_response(response: &MessageResponse) -> bool { + match &response.message { + Message::Assistant { content } => { + content.iter().any(|b| matches!(b, ContentBlock::ToolUse { .. })) + } + _ => false, + } +} + +fn extract_tool_calls_from_response(response: &MessageResponse) -> Vec<(String, String, String)> { + match &response.message { + Message::Assistant { content } => { + content.iter().filter_map(|b| { + if let ContentBlock::ToolUse { id, name, input } = b { + // input 是 &Value,Display 实现生成紧凑 JSON + let args = input.to_string(); + Some((id.clone(), name.clone(), args)) + } else { + None + } + }).collect() + } + _ => vec![], + } +} +``` + +### 4b:工具结果回传(依赖 FIX-A) + +> ⚠️ **前置依赖**:`ToolInvocation` 必须在 Phase 0/1 期间增加 `tool_call_id` 字段(见 FIX-A),否则以下代码无法编译。 + +```rust +for result in results { + let is_error = result.output.is_err(); + let content = match result.output { + Ok(ref value) => { + let serialized = serde_json::to_string(value).unwrap_or_default(); + truncate_tool_result(&serialized, max_bytes) + } + Err(ref e) if e.is_recoverable() => format!("错误: {}", e), + Err(ref e) => return Err(LlmError::Other(format!("工具不可恢复错误: {}", e))), + }; + self.messages.push(Message::tool_result( + result.tool_call_id.clone(), // ← 来自 FIX-A 新增字段 + content, + is_error, + )); +} +``` + +**数据流**(FIX-A 修复后): +1. `extract_tool_calls_from_response()` 返回 `Vec<(String, String, String)>`(id, name, args) +2. `invoke_all()` 接收 `Vec<(String, String, Value)>`(id, name, args) +3. `ToolInvocation.tool_call_id` 保存原始 id +4. 回传时使用 `result.tool_call_id` 关联原始调用 + +--- + +## 任务 5:简化 `submit()` / `submit_messages()` / `submit_request()` + +### 5a:`submit()` + +```rust +pub async fn submit( + &mut self, + prompt: String, + tools: Vec, +) -> Result { + self.messages.push(Message::user_text(prompt)); + self.maybe_compact(); + + let mut attempts = 0; + loop { + let request = self.build_request(&tools); + // ... hook 处理 ... + match self.provider.chat(request).await { + Ok(response) => { + // ... hook 处理 ... + self.messages.push(response.message.clone()); + self.usage.add(&response.usage); + return Ok(response); + } + Err(e) => { /* 重试逻辑 */ } + } + } +} +``` + +**变更**: +- `self.messages.push(OpenaiChatMessage::user_text(prompt))` → `self.messages.push(Message::user_text(prompt))` +- `self.messages.push(response.message.clone())` — `response.message` 已经是 `Message` 类型 +- 无需 `message_to_chat_message` 转换 + +### 5b:`submit_messages()` + +```rust +pub async fn submit_messages( + &mut self, + messages: Vec, + tools: Vec, +) -> Result { + let request = MessageRequest { + model: self.config.model.clone(), + messages, + tools: tools.to_vec(), + tool_choice: ToolChoice::Auto, + max_tokens: self.config.max_tokens, + temperature: self.config.temperature, + ..Default::default() + }; + // ... hook + provider.chat(request) ... + match self.provider.chat(request).await { + Ok(response) => { + self.usage.add(&response.usage); + Ok(response) + } + Err(e) => Err(e), + } +} +``` + +### 5c:`submit_request()` + +```rust +async fn submit_request( + &mut self, + tools: &[ToolDefinition], +) -> Result { + // 逻辑与之前一致,但返回类型改为 MessageResponse + // ... +} +``` + +--- + +## 任务 6:适配 `compact.rs` + +### 6a:`estimate_single_message_tokens()` + +```rust +fn estimate_single_message_tokens(msg: &Message) -> u32 { + let role_overhead: u32 = 4; + let content_tokens = match msg { + Message::System { content } + | Message::User { content } + | Message::Assistant { content } => estimate_content_blocks_tokens(content), + Message::UserImage { .. } => { + // UserImage 固定估算 + 50 + } + Message::ToolResult { content, .. } => estimate_content_blocks_tokens(content), + }; + role_overhead + content_tokens +} +``` + +**变更**:当前使用 `ContentField`(旧类型,区分 String 和 Array),新类型统一为 `&[ContentBlock]`。 + +### 6b:`estimate_content_blocks_tokens()` + +```rust +fn estimate_content_blocks_tokens(blocks: &[ContentBlock]) -> u32 { + blocks.iter().map(estimate_block_tokens).sum() +} + +fn estimate_block_tokens(block: &ContentBlock) -> u32 { + match block { + ContentBlock::Text { text } => estimate_text_tokens(text), + ContentBlock::ToolResult { content, .. } => estimate_content_blocks_tokens(content), + ContentBlock::Thinking { text, .. } => estimate_text_tokens(text), + ContentBlock::ToolUse { input, .. } => { + // tool_use 的 input 是 JSON Value + estimate_text_tokens(&serde_json::to_string(input).unwrap_or_default()) + } + _ => 50, // Image/Audio/File/Extension 固定估算 + } +} +``` + +### 6c:`microcompact()`(审查 FIX-F) + +> ⚠️ **FIX-F**:只压缩 `is_error: false` 的 ToolResult。错误结果包含失败原因,压缩后 LLM 无法理解。 + +```rust +pub fn microcompact(messages: &mut [Message], keep_recent: usize) -> u32 { + if messages.len() <= keep_recent { + return 0; + } + let prune_start = messages.len() - keep_recent; + let mut freed_tokens: u32 = 0; + + // 第一遍:计算可释放的 token(仅非错误结果) + for msg in &messages[..prune_start] { + if matches!(msg, Message::ToolResult { is_error: false, .. }) { + freed_tokens += estimate_single_message_tokens(msg); + } + } + + // 第二遍:替换内容(仅非错误结果) + for msg in &mut messages[..prune_start] { + if let Message::ToolResult { content, is_error: false } = msg { + *content = vec![ContentBlock::Text { text: "[pruned]".to_string() }]; + } + } + + freed_tokens +} +``` + +**变更**: +- `OpenaiChatMessage::Tool` → `Message::ToolResult` +- `ContentField::Array(vec![OpenaiContentPart::Text { text: "[pruned]" }])` → `vec![ContentBlock::Text { text: "[pruned]".to_string() }]` +- 新增条件:`is_error: false` 时才压缩(保留错误诊断信息) + +### 6d:`should_compact()` + +```rust +pub fn should_compact( + messages: &[Message], + config: &CompactConfig, + state: &CompactState, +) -> bool { + // 逻辑不变,只是消息类型改为 &[Message] + if state.consecutive_failures >= MAX_CONSECUTIVE_FAILURES { + return false; + } + let tokens = estimate_message_tokens(messages); + tokens >= config.threshold() +} +``` + +--- + +## 任务 7:同步适配 `ConversationMemory` + +`src/memory/conversation.rs` 的消息类型从 `OpenaiChatMessage` 切换为 `Message`: + +### 7a:结构体字段 + +```rust +pub struct ConversationMemory { + store: Arc, + session_id: String, + config: ConversationMemoryConfig, + messages: Vec, // 从 OpenaiChatMessage 改为 Message + message_ids: Vec, + compact_state: CompactState, +} +``` + +### 7b:序列化/反序列化 + +```rust +// 当前 +serde_json::from_str::(&item.content) + +// 目标 +serde_json::from_str::(&item.content) +``` + +**注意**:`Message` 枚举需要实现 `Serialize` / `Deserialize`。 + +> ✅ **已在 Phase 0 处理**(审查 FIX-B):Phase 0 定义 `Message` 和 `ContentBlock` 类型时直接添加 `#[derive(Serialize, Deserialize)]` 和 `#[serde(tag = "type", rename_all = "snake_case")]`。Phase 2 直接使用即可。 + +序列化格式示例(仅用于内部存储,不与 Provider 协议交互): + +```json +{"type": "user", "content": [{"type": "text", "text": "hello"}]} +{"type": "user_image", "data": "...", "mime_type": "image/png", "detail": "auto"} +{"type": "assistant", "content": [{"type": "text", "text": "Hi!"}, {"type": "tool_use", "id": "call_1", "name": "get_weather", "input": {"city": "Beijing"}}]} +{"type": "tool_result", "tool_call_id": "call_1", "content": [{"type": "text", "text": "Sunny, 25°C"}], "is_error": false} +``` + +### 7c:消息 ID 前缀兼容 + +```rust +fn session_prefix(&self) -> String { + format!("conv:{}:", self.session_id) // 格式不变 +} +``` + +序列化格式变化后,旧消息无法反序列化为新类型。由于项目尚未 release,可以接受这个不兼容。 + +--- + +## 任务 8:清理临时转换函数 + +删除 Phase 0 引入的以下函数(确认不再被引用后删除): + +1. **`chat_message_to_message()`** — 在 `build_request()` 中已不再需要(`self.messages` 已是 `Vec`) +2. **`message_to_chat_message()`** — 在 `submit_stream()` 中已不再需要(流结束后直接使用 `full_response.message`) +3. **`content_to_blocks()`** — 如果只在上述两个函数中使用,一并删除 +4. **`blocks_to_content_field()`** — 同上 +5. **`has_tool_calls_in_message()`** — `cycle.rs:645`,基于 `OpenaiChatMessage::Assistant { tool_calls }` 判断。Phase 2 中由 `has_tool_calls_in_response()`(基于 `ContentBlock::ToolUse`)替代 +6. **`extract_tool_calls_from_message()`** — `cycle.rs:658`,基于 `OpenaiToolCall` 提取。Phase 2 中由 `extract_tool_calls_from_response()`(基于 `ContentBlock::ToolUse`)替代 + +**验证方法**: +- 全局搜索函数名(`git grep`),确认在所有文件中的引用计数归零 +- `cargo build` 通过(编译器对非 pub 的未使用函数报 warning,pub 函数不报 warning。用 `git grep` 做精准确认) +- `cargo clippy` 无 `dead_code` 级别警告 + +--- + +## 任务 9:适配 `prompt/composer.rs`(如需要) + +`validate_messages()` 当前检查 `OpenaiChatMessage` 序列中的 role 顺序规则。如果该函数被 `AgentSession` 或 `LlmCycle` 调用,需要适配新 `Message` 类型。 + +检查步骤: +1. 搜索 `validate_messages` 的调用者 +2. 如果被上游调用且上游已切换为 `Vec`,则适配 +3. 适配内容:match 分支从 `OpenaiChatMessage::Tool { .. }` 改为 `Message::ToolResult { .. }`,从 `OpenaiChatMessage::Assistant { tool_calls: Some(..) }` 改为检查 `ContentBlock::ToolUse` + +--- + +## 验证方式 + +1. **编译检查**:`cargo build` 通过 +2. **单元测试**:`cargo test` 全部通过 +3. **集成测试**: + - 多轮对话(`submit` + `submit_with_tools`)端到端正常 + - 工具调用循环正常(含 `tool_call_id` 正确关联) + - 流式响应(`submit_stream`)正常 + - 上下文压缩(`compact`)超过 token 阈值后正确压缩 + - **错误路径**:`provider.chat()` 返回错误时,消息历史不被污染(`self.messages.push` 在确认成功后才执行) + - **序列化 roundtrip**:`Message` → JSON → `Message` 保持所有字段(含 `UserImage`、`ToolResult.is_error`) + - **边界情况**:`submit_with_tools` 收到空 content 的 `ContentBlock`、包含 `UserImage` 的消息 + - **hook 集成**:PreRequest hook 正确收到 `MessageRequest`(而非旧的 `ChatRequest`) +4. **clippy 检查**:`cargo clippy` 无新增警告 +5. **diff 确认**:`git diff` 确认 Phase 0 引入的临时转换函数已被删除 + +## 回滚方案 + +⚠️ **审查评估**:feature flag 方案在两套 StreamEvent 类型共存时工程量大,实际安全网依赖 checkpoint 打 tag。 + +| 回退对象 | 可行策略 | 复杂度 | 说明 | +|---------|---------|-------|------| +| `submit_stream()` | feature flag `legacy_stream` | 🔴 高 | 需要同时保留旧 `StreamEvent` 类型和旧 `chat_stream` Provider 调用路径。除非新旧 StreamEvent 完全兼容,否则工程量大。**不推荐** | +| `compact.rs` | 模块同步回退 | 🟢 低 | `should_compact`/`microcompact` 只在 `LlmCycle` 和 `ConversationMemory` 中调用,签名变化后可整体回退 | +| 全局 | git checkout + 重做 | 🟢 极低 | 项目尚未 release,无外部消费者。每个子阶段完成后打 tag(`phase2-taskX-checkpoint`),允许跳跃回退 | + +**推荐回退策略**:不在代码层面做 feature flag,而是: +1. 每个子阶段(Phase 2.0/2.1/2.2/2.3)完成后创建 git tag +2. 如果某一子阶段出现问题,`git revert` 该子阶段的所有 commit +3. 修复后重新提交,不保留两套实现共存 + +## 开放事项 + +- [已解决] `Message` 的 `Serialize` / `Deserialize` 派生设计 → **移至 Phase 0**(FIX-B) +- [已解决] `system_prompt` 字段的最终去留 → **Phase 2 彻底移除**(FIX-D) +- `ConversationMemory` 序列化格式变更后的数据兼容性(项目未 release,可接受) +- `submit_with_tools_stream()` 未来设计:流式响应 + 自动工具循环的组合场景 +- `ToolContext` 是否增加 `tool_call_id` 字段(FIX-A 方案二中可选) +- `ContentBlock::Extension` 逃生舱的具体使用场景 +- Thinking signature 的端到端测试