docs(9): 完善 Anthropic 流式状态机设计推演结论
This commit is contained in:
@@ -105,7 +105,7 @@ pub type Message = OpenaiChatMessage;
|
||||
| 3 | StreamEvent 汇聚为 MessageResponse 算法 | [9c-llm-provider-trait.md](9c-llm-provider-trait.md#44-partialmessageresponse--流式事件的汇聚算法) | ✅ 已推演(方案 B:显式边界 + BTreeMap 分桶) |
|
||||
| 4 | ToolCallStart index 归一化 | [9c-llm-provider-trait.md](9c-llm-provider-trait.md#43-流式事件-streamevent) | ✅ 已推演(自动解决,ToolCallStart 合并到 ContentBlockStart) |
|
||||
| 5 | OpenAI 流式转换实现 | [9d-provider-implementations.md](9d-provider-implementations.md#51-openai-provider兼容-chat-api) | ✅ 已推演(方案 A:SseByteStream 通用层 + OpenaiStreamToEvents 状态机,ToolCallEnd 依赖 finalize 兜底,忽略多 Choice) |
|
||||
| 6 | Anthropic 流式状态机设计 | [9d-provider-implementations.md](9d-provider-implementations.md#52-anthropicprovidermessages-api) | **高** |
|
||||
| 6 | Anthropic 流式状态机设计 | [9d-provider-implementations.md](9d-provider-implementations.md#52-anthropicprovidermessages-api) | ✅ 已推演(轻量分发器:3 状态 + 7 种事件映射 + 零 index 映射) |
|
||||
| 7 | OpenAI Response API 完整映射表 | [9d-provider-implementations.md](9d-provider-implementations.md#53-openai-response-api草案) | 低 |
|
||||
| 8 | DeepSeek/Qwen Provider 落地策略 | [9d-provider-implementations.md](9d-provider-implementations.md#54-deepseek--qwen-等兼容-provider-的落地策略) | 低 |
|
||||
| 9 | system prompt 双重表达冲突 | [9e-llm-cycle-and-upstream.md](9e-llm-cycle-and-upstream.md#62-build_request--新签名) | **高** |
|
||||
|
||||
@@ -70,25 +70,31 @@ OpenaiChatResponse → MessageResponse:
|
||||
> │
|
||||
> ▼
|
||||
> SseByteStream [sse.rs — 通用层]
|
||||
> │ 逐行分割、去 "data:" 前缀、过滤 "[DONE]"、缓冲拼接碎片
|
||||
> │ 逐行分割、按空行分帧、解析 event:/data: 行前缀、
|
||||
> │ 过滤 "[DONE]"/ping、缓冲拼接碎片、多行 data: 自动拼接
|
||||
> ▼
|
||||
> Stream<Item=Result<String, LlmError>> ← JSON 字符串行
|
||||
> Stream<Item=Result<SseEvent, LlmError>> ← 结构化 SSE 帧
|
||||
> │ SseEvent { event_name: Option<String>, data: String }
|
||||
> │ - OpenAI: event_name = None(未命名事件)
|
||||
> │ - Anthropic: event_name = Some("message_start" | "content_block_delta" | ...)
|
||||
> │
|
||||
> ▼
|
||||
> OpenaiStreamToEvents [openai/stream.rs — OpenAI 特定]
|
||||
> │ 反序列化为 OpenaiChatChunk → 状态机转换
|
||||
> ▼
|
||||
> Stream<Item=Result<StreamEvent, LlmError>> ← IR 语义事件
|
||||
> ├─ OpenAI ────→ OpenaiStreamToEvents [openai/stream.rs]
|
||||
> │ 忽略 event_name,反序列化 data → OpenaiChatChunk → 状态机
|
||||
> │
|
||||
> └─ Anthropic ──→ AnthropicStreamToEvents [anthropic/stream.rs]
|
||||
> 按 event_name 分发事件类型 → 事件映射
|
||||
> ```
|
||||
>
|
||||
> **Layer 1 — `SseByteStream<S>`(新增 `src/llm/provider/sse.rs`)**
|
||||
>
|
||||
> 从当前 `SseChunkStream`(`openai.rs` 内联实现)提取字节解析逻辑为通用 SSE 解析器。
|
||||
> 产出 JSON 字符串行,不绑定 OpenAI 格式。Anthropic 直接复用。
|
||||
> 产出 `SseEvent` 结构体,携带 `event_name` + `data` 两部分信息。
|
||||
> OpenAI 和 Anthropic 均直接复用此层,各自的事件转换器按需使用 `event_name`。
|
||||
>
|
||||
> **Layer 2 — `OpenaiStreamToEvents`(新增 `src/llm/provider/openai/stream.rs`)**
|
||||
>
|
||||
> 接收 JSON 字符串行,反序列化为 `OpenaiChatChunk`(保留作为内部格式),
|
||||
> 接收 `SseEvent` 流,忽略 `event_name`(OpenAI 为 `None`),
|
||||
> 将 `data` 反序列化为 `OpenaiChatChunk`(保留作为内部格式),
|
||||
> 通过状态机转换为 `StreamEvent` 事件流。
|
||||
>
|
||||
> ### 状态机设计
|
||||
@@ -143,7 +149,7 @@ OpenaiChatResponse → MessageResponse:
|
||||
>
|
||||
> | # | 议题 | 结论 | 理由 |
|
||||
> |---|------|------|------|
|
||||
> | 1 | **SSE 字节解析复用** | 提取为通用 `SseByteStream`(`sse.rs`),`OpenaiStreamToEvents` 在其上封装 | Anthropic 复用相同字节协议、分层测试、职责清晰 |
|
||||
> | 1 | **SSE 字节解析复用** | 提取为通用 `SseByteStream`(`sse.rs`),支持 `event:` + `data:` 双行解析;`OpenaiStreamToEvents` / `AnthropicStreamToEvents` 分别在其上封装 | Anthropic 复用相同字节协议,通过 `SseEvent.event_name` 区分事件类型;分层测试、职责清晰 |
|
||||
> | 2 | **ContentBlockStart/End 合成** | "类型切换推断边界"策略——来什么类型就关旧开新,依赖 content/tool_calls 互斥保证 | 状态机 3 种当前类型覆盖全部场景,假设验证通过 |
|
||||
> | 3 | **Tool call 序号映射** | `HashMap<openai_index, global_block_index>`,新 tool call 出现时分配全局序号 | OpenAI index 是 tool 数组级别全局的,但与 IR content block 体系不同,需映射 |
|
||||
> | 4 | **多 Choice** | 忽略 choices[1..],不暴露 | 当前架构无多 choice 概念,80% 场景 n=1,非目标已明确 |
|
||||
@@ -168,7 +174,7 @@ OpenaiChatResponse → MessageResponse:
|
||||
>
|
||||
> | 操作 | 文件 | 说明 |
|
||||
> |------|------|------|
|
||||
> | 新增 | `src/llm/provider/sse.rs` | 通用 SSE 字节解析层,从当前 `openai.rs` 的 `SseChunkStream` 提取核心逻辑 |
|
||||
> | 新增 | `src/llm/provider/sse.rs` | 通用 SSE 字节解析层,含 `SseEvent` 结构体;从当前 `openai.rs` 的 `SseChunkStream` 提取核心逻辑并增强为支持 `event:` + `data:` 双行解析 + 空行分帧 |
|
||||
> | 新增 | `src/llm/provider/openai/stream.rs` | `OpenaiStreamToEvents` 转换器 + 状态机 |
|
||||
> | 修改 | `src/llm/provider/openai.rs` | `chat_stream` 返回 `Result<StreamEvent>`,组合 `SseByteStream` + `OpenaiStreamToEvents` |
|
||||
> | 修改 | `src/llm/provider.rs` | `LlmProvider::chat_stream` 签名改为 `Result<StreamEvent>` |
|
||||
@@ -216,39 +222,181 @@ Anthropic Response → MessageResponse:
|
||||
stop_reason → stop_reason
|
||||
```
|
||||
|
||||
> **🔄 待深入推演:Anthropic 流式状态机设计**
|
||||
> Anthropic 的流式比 OpenAI 复杂得多——7 种事件类型、需要维护 block index 状态、
|
||||
> thinking 有 signature 晚于 content_block 下发。
|
||||
> **需要推演:**
|
||||
> 1. **状态机状态设计**:
|
||||
> - `PendingStart` → 等待 `message_start`
|
||||
> - `InBlock { block_index, block_type, tool_acc: Option<ToolAccumulator> }` → 在某个 content block 中
|
||||
> - `BetweenBlocks` → 等待下一个 `content_block_start` 或 `message_delta`
|
||||
> - `Completed` → 收到 `message_stop`
|
||||
> - `Errored`
|
||||
> 2. **block index 追踪**:`content_block_start` 中的 `index` 是全局 content block 序号,
|
||||
> 直接对应最终 content 数组中的位置。`ToolUse { id, name }` 嵌入在 `ContentBlockStart` 的
|
||||
> `block_type` 中(已推演:ToolCallStart 已移除,合并到 ContentBlockType::ToolUse)。
|
||||
> 而 `input` 通过后续的 `content_block_delta`(含 `input_json_delta`)增量到达。
|
||||
> 实现时发出 `ContentBlockStart(index, ToolUse { id, name })` 后,
|
||||
> 通过后续的 `ToolCallArgumentsDelta { index, arguments }` 累积参数。
|
||||
> 3. **Thinking signature 的附着时机**:**✅ 已推演(方案 C:MessageComplete 兜底)**
|
||||
> Anthropic 中,thinking block 的 `signature` 不在 `content_block_start` 或 `content_block_delta`
|
||||
> 中下发,而是在最后的 `message_delta`(与 stop_reason 一起)下发。
|
||||
> **推演结论:**
|
||||
> - `StreamEvent::MessageComplete { stop_reason, thinking_signature: Option<String> }`
|
||||
> 携带 signature,不在 ThinkingDelta 或 ContentBlockEnd 中传递
|
||||
> - Anthropic 映射层在收到 `message_delta` 时,将 `message_delta.thinking.signature`
|
||||
> 填入 `MessageComplete.thinking_signature`
|
||||
> - 汇聚算法(`PartialMessageResponse::finalize()`)在遍历 Thinking block 时,
|
||||
> 如果 builder 中的 `signature` 为 None,用 `self.thinking_signature` 回填
|
||||
> - 该方案不需新增独立事件,统一在 finalize 时处理,语义清晰
|
||||
> 4. **block 嵌套的边界情况**:Anthropic 的响应中 block 不会嵌套,但允许多个 block 连续出现。
|
||||
> 需要确保状态机正确处理 block 间切换(前一个 `content_block_stop` → 后一个 `content_block_start`)。
|
||||
> **需要推演:**
|
||||
> - 完整的状态转移图(当前状态 + 输入事件 → 新状态 + 产出的 StreamEvent)
|
||||
> - SSE 字节流解析器(与 OpenAI 共享字节流层,差异化事件解析层)
|
||||
> - 错误恢复:收到 `error` 事件时的状态机行为
|
||||
> **✅ 推演结论(2026-06-18):**
|
||||
>
|
||||
> ### 设计思路:轻量分发器
|
||||
>
|
||||
> Anthropic 的流式事件**自带语义块边界**(`content_block_start/stop` 显式声明),
|
||||
> 不像 OpenAI 需从扁平 delta 推断。因此 Anthropic 状态机采用**轻量分发器**模式——
|
||||
> 每个事件自描述,状态机仅做顺序合法性校验,不做 block 边界推断或 index 映射。
|
||||
>
|
||||
> **与 OpenAI 流式转换的核心差异:**
|
||||
>
|
||||
> | 维度 | OpenaiStreamToEvents | AnthropicStreamToEvents |
|
||||
> |------|---------------------|------------------------|
|
||||
> | 核心复杂度 | 中——需从扁平 delta 推断 block 边界 | 低——事件自带语义边界 |
|
||||
> | 状态数 | 3(无活跃、Text、ToolUse) | 3(PendingStart、Active、Terminated) |
|
||||
> | index 管理 | `HashMap<openai_idx, global_idx>` 映射 | 直接使用 Anthropic index(1:1) |
|
||||
> | Block 边界推断 | 类型切换推断 | 原生 content_block_start/stop |
|
||||
> | Thinking 处理 | 无 | 通过 message_delta.thinking.signature |
|
||||
>
|
||||
> ### 架构分层
|
||||
>
|
||||
> 沿用 OpenAI 的两层架构,`SseByteStream` 共享(已增强为支持命名事件),
|
||||
> `AnthropicStreamToEvents` 在事件层按 `event_name` 分发:
|
||||
>
|
||||
> ```
|
||||
> bytes_stream()
|
||||
> │
|
||||
> ▼
|
||||
> SseByteStream [sse.rs — 通用层(已增强)]
|
||||
> │ 逐行分割、按空行分帧、解析 event:/data: 行前缀、过滤 "[DONE]"/ping
|
||||
> ▼
|
||||
> Stream<Item=Result<SseEvent, LlmError>> ← SseEvent { event_name: Option<String>, data }
|
||||
> │
|
||||
> ▼
|
||||
> AnthropicStreamToEvents [anthropic/stream.rs — Anthropic 特定]
|
||||
> │ 按 SseEvent.event_name 分发事件类型 → 直接映射
|
||||
> ▼
|
||||
> Stream<Item=Result<StreamEvent, LlmError>> ← IR 语义事件
|
||||
> ```
|
||||
>
|
||||
> ### 结构体设计
|
||||
>
|
||||
> ```rust
|
||||
> pub struct AnthropicStreamToEvents<S> {
|
||||
> inner: S, // Stream<Item=Result<SseEvent, LlmError>>
|
||||
> state: AnthropicStreamState, // 仅做顺序校验
|
||||
> usage: PartialUsage, // 从 message_start + message_delta 累积
|
||||
> pending_thinking_signature: Option<String>, // message_delta 中到达
|
||||
> }
|
||||
>
|
||||
> /// 状态机状态 —— 仅做合法性校验,事件本身已自描述。
|
||||
> enum AnthropicStreamState {
|
||||
> PendingStart, // 等待 message_start
|
||||
> Active, // 已收到 message_start,正在接收 content block 事件
|
||||
> Terminated, // 已终结,不再处理后续事件
|
||||
> }
|
||||
> ```
|
||||
>
|
||||
> 状态足够简单的原因:Anthropic 每个事件自带完整语义——
|
||||
> - `content_block_delta` 自带 `index`,不需要追踪"当前活跃 block"
|
||||
> - `content_block_stop` 自带 `index`,不需要追踪"当前关闭哪个"
|
||||
> - 状态只拒绝非法到达顺序的事件
|
||||
>
|
||||
> ### 事件映射表(完整)
|
||||
>
|
||||
> | Anthropic SSE 事件 | 产出的 StreamEvent | 说明 |
|
||||
> |-------------------|-------------------|------|
|
||||
> | `message_start` | `MessageStart { id, model }`<br>`CostUpdate { prompt_tokens }` | 从 `message.usage.input_tokens` 提取 |
|
||||
> | `ping` | —(忽略) | Anthropic 心跳 |
|
||||
> | `content_block_start`<br>`block.type="text"` | `ContentBlockStart { index, Text }` | |
|
||||
> | `content_block_start`<br>`block.type="tool_use"` | `ContentBlockStart { index, ToolUse { id, name } }` | block 自带 id + name |
|
||||
> | `content_block_start`<br>`block.type="thinking"` | `ContentBlockStart { index, Thinking }` | |
|
||||
> | `content_block_delta`<br>`delta.type="text_delta"` | `TextDelta { text: delta.text }` | |
|
||||
> | `content_block_delta`<br>`delta.type="thinking_delta"` | `ThinkingDelta { text: delta.thinking }` | |
|
||||
> | `content_block_delta`<br>`delta.type="input_json_delta"` | `ToolCallArgumentsDelta { index, arguments: delta.partial_json }` | index 透传 |
|
||||
> | `content_block_stop` | `ContentBlockEnd { index }` | |
|
||||
> | `message_delta` | `CostUpdate { completion_tokens }`<br>`MessageComplete { stop_reason, thinking_signature }` | signature 从 `delta.thinking?.signature` 提取 |
|
||||
> | `message_stop` | —(流结束标记,不产事件) | 仅切状态到 Terminated |
|
||||
> | `error` | `Error { message: error.message }` | 切状态到 Terminated |
|
||||
>
|
||||
> ### 状态转移表
|
||||
>
|
||||
> **当前状态:`PendingStart`**
|
||||
>
|
||||
> | 输入事件 | 输出 StreamEvent | 新状态 | 备注 |
|
||||
> |---------|----------------|--------|------|
|
||||
> | `message_start` | → `MessageStart` + `CostUpdate` | `Active` | ✅ 正常流程入口 |
|
||||
> | 其他任何事件 | — ⚠ warn | 不变 | 防御性跳过 |
|
||||
> | `error` | → `Error` | `Terminated` | ❌ 错误路径 |
|
||||
>
|
||||
> **当前状态:`Active`**
|
||||
>
|
||||
> | 输入事件 | 输出 StreamEvent | 新状态 | 备注 |
|
||||
> |---------|----------------|--------|------|
|
||||
> | `ping` | —(忽略) | `Active` | ✅ 心跳 |
|
||||
> | `content_block_start` | → `ContentBlockStart` | `Active` | ✅ 新 block 开始 |
|
||||
> | `content_block_delta` | → `TextDelta` / `ThinkingDelta` / `ToolCallArgumentsDelta` | `Active` | ✅ 块内增量 |
|
||||
> | `content_block_stop` | → `ContentBlockEnd` | `Active` | ✅ block 结束 |
|
||||
> | `message_delta` | → `CostUpdate` + `MessageComplete` | `Active` | ✅ 消息完成信息 |
|
||||
> | `message_stop` | — | `Terminated` | ✅ 正常结束 |
|
||||
> | `error` | → `Error` | `Terminated` | ❌ 错误路径 |
|
||||
> | 未知 delta type | — ⚠ warn | `Active` | 防御性忽略 |
|
||||
>
|
||||
> **当前状态:`Terminated`**
|
||||
>
|
||||
> | 输入事件 | 输出 StreamEvent | 新状态 | 备注 |
|
||||
> |---------|----------------|--------|------|
|
||||
> | 任何事件 | — ⚠ warn "已完结" | `Terminated` | 防御性忽略 |
|
||||
>
|
||||
> ### 各议题结论
|
||||
>
|
||||
> | # | 议题 | 结论 | 理由 |
|
||||
> |---|------|------|------|
|
||||
> | 1 | **状态机模式** | 轻量分发器(3 状态),不做 block 推断 | Anthropic 事件自带语义边界,不需要像 OpenAI 那样推断 |
|
||||
> | 2 | **index 映射** | 无需映射,直接使用 Anthropic index(1:1) | Anthropic 的 `index` 是全局 content block 序号,与 IR 完全对齐 |
|
||||
> | 3 | **Thinking signature** | ✅ 已推演(方案 C):message_delta 提取 → MessageComplete 传递 → finalize 回填 | 已在 [9f-edge-cases.md](9f-edge-cases.md#92-thinking-的端到端流程) 中完成推演 |
|
||||
> | 4 | **SSE 字节解析复用** | 通过增强的 `SseByteStream`(支持 `event:` 行解析)与 OpenAI 共享通用层 | 同一字节协议,仅在事件解析层差异化 |
|
||||
> | 5 | **Usage 分次到达** | `message_start` 提取 `input_tokens`,`message_delta` 提取 `output_tokens`,`PartialUsage` 字段级合并 | 与 StreamEvent 汇聚算法兼容 |
|
||||
> | 6 | **错误恢复** | `error` 事件 → `StreamEvent::Error` + state=Terminated,后续事件全部忽略 | 不同于 OpenAI 的 HTTP 错误路径,但 IR 层统一为 `StreamEvent::Error` |
|
||||
> | 7 | **ContentBlockType::ToolUse 嵌入** | `content_block_start` 中的 `id` + `name` 直接填入 `ContentBlockType::ToolUse { id, name }` | 与已推演的 ContentBlockStart 设计一致 |
|
||||
> | 8 | **block 切换** | content_block_stop(index) → content_block_start(index') 自然过渡,状态机不追踪 | 事件本身已确定边界,无需状态机参与 |
|
||||
>
|
||||
> ### 与已推演设计的对齐
|
||||
>
|
||||
> **与 Thinking signature 方案的对齐(方案 C):**
|
||||
> ```
|
||||
> content_block_start { type: "thinking" }
|
||||
> → ContentBlockStart(Thinking) ← signature 未到达
|
||||
> content_block_delta { thinking_delta }
|
||||
> → ThinkingDelta(...)
|
||||
> content_block_stop
|
||||
> → ContentBlockEnd ← signature 仍未到达
|
||||
> message_delta { delta.thinking.signature = "0x..." }
|
||||
> → MessageComplete { thinking_signature: Some("0x...") }
|
||||
> → finalize() 回填到最后一个 Thinking block
|
||||
> ```
|
||||
>
|
||||
> **与 StreamEvent 汇聚算法的对齐:**
|
||||
> 本状态机产出的 StreamEvent 可直接喂入已推演的 `PartialMessageResponse::apply_to()` 算法。
|
||||
> 上述事件序列在汇聚算法中:
|
||||
> 1. `MessageStart` → state.id/model
|
||||
> 2. `ContentBlockStart/Delta/End` → `BTreeMap` 按 index 分桶组装
|
||||
> 3. `CostUpdate` → PartialUsage 字段级合并
|
||||
> 4. `MessageComplete` → stop_reason + thinking_signature + is_complete
|
||||
> 5. `finalize()` → 回填 signature → `MessageResponse`
|
||||
>
|
||||
> ### 边界情况
|
||||
>
|
||||
> | 场景 | 处理方式 |
|
||||
> |------|---------|
|
||||
> | **message_start 前收 content_block_start** | ⚠ warn 忽略,不发射事件 |
|
||||
> | **message_delta 前收 message_stop** | ⚠ warn,强制 Terminated |
|
||||
> | **content_block_stop 无对应 start** | ⚠ warn 忽略(index 无对应) |
|
||||
> | **index 跳跃(0 → 2)** | 正常处理,index 透传,content 数组留空位 |
|
||||
> | **delta index 与最新 start 不匹配** | ⚠ warn,仍然按 delta 自带 index 处理 |
|
||||
> | **message_delta 缺 thinking.signature** | `thinking_signature`: None |
|
||||
> | **message_delta 缺 usage** | 不发射 CostUpdate,仅发射 MessageComplete |
|
||||
> | **两次 message_delta** | ⚠ warn,第二次忽略 |
|
||||
> | **content_block_stop 后同 index 又来 delta** | ⚠ warn 忽略 |
|
||||
> | **ping 事件** | 忽略,不发射任何事件 |
|
||||
> | **网络断开** | emit `Error`,state = Terminated |
|
||||
> | **JSON 解析失败** | emit `Error`,state = Terminated |
|
||||
> | **stop_reason 映射** | `"end_turn"`→`Stop`, `"max_tokens"`→`MaxTokens`, `"tool_use"`→`ToolUse`, `"stop_sequence"`→`StopSequence`, 其他→`Other` |
|
||||
>
|
||||
> ### 文件变更清单
|
||||
>
|
||||
> | 操作 | 文件 | 说明 |
|
||||
> |------|------|------|
|
||||
> | 新增 | `src/llm/provider/anthropic.rs` | AnthropicProvider 实现(chat + chat_stream) |
|
||||
> | 新增 | `src/llm/provider/anthropic/` | 目录,按 2018 版风格组织 |
|
||||
> | 新增 | `src/llm/provider/anthropic/stream.rs` | `AnthropicStreamToEvents` 转换器 + 事件分发器 |
|
||||
> | 增强 | `src/llm/provider/sse.rs` | `SseByteStream` 增强为支持 `event:` 行 + 空行分帧(已在 §5.1 中描述) |
|
||||
> | 修改 | `src/llm/provider.rs` | `ProviderType` 增加 `Anthropic`;`create_provider` 增加分支 |
|
||||
> | 无变更 | `src/llm/cycle.rs` | StreamEvent 事件序列格式不变,无需改动 |
|
||||
> | 无变更 | `src/llm/stream.rs` | 汇聚算法 `apply_to/finalize` 不变 |
|
||||
>
|
||||
> 优先级:高(Phase 4 AnthropicProvider 实现的前提条件)
|
||||
|
||||
### 5.3 OpenAI Response API(草案)
|
||||
|
||||
Reference in New Issue
Block a user