docs(9d): 补全 OpenAI 流式转换推演结论

This commit is contained in:
徐涛
2026-06-18 06:34:20 +08:00
parent c0eae92b10
commit 8c9324350b
2 changed files with 117 additions and 22 deletions
+1 -1
View File
@@ -104,7 +104,7 @@ pub type Message = OpenaiChatMessage;
| 2 | extra 的类型安全性 | [9b-ir-type-system.md](9b-ir-type-system.md#36-messagerequest--统一请求) | ✅ 已推演(方案 BResult-based access |
| 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) | **高** |
| 5 | OpenAI 流式转换实现 | [9d-provider-implementations.md](9d-provider-implementations.md#51-openai-provider兼容-chat-api) | ✅ 已推演(方案 ASseByteStream 通用层 + OpenaiStreamToEvents 状态机,ToolCallEnd 依赖 finalize 兜底,忽略多 Choice |
| 6 | Anthropic 流式状态机设计 | [9d-provider-implementations.md](9d-provider-implementations.md#52-anthropicprovidermessages-api) | **高** |
| 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-的落地策略) | 低 |
+116 -21
View File
@@ -59,27 +59,122 @@ OpenaiChatResponse → MessageResponse:
}
```
> **🔄 待深入推演:OpenAI 流式转换的具体实现**
> 当前只有一行从 OpenAI SSE chunk 到 StreamEvent 的映射,没有实现细节。
> 实际实现中需要解决以下问题:
> 1. **SSE 字节流解析器**:当前 `SseChunkStream` 从字节流解析 SSE 行(`data: ...`),新设计中需要改写为
> 直接产出 StreamEvent 的转换器。是否可以复用现有 `SseChunkStream` 的字节流解析逻辑?
> 2. **ContentBlockStart/End 合成**OpenAI 的 SSE 流中没有原生 block 边界标记,只有一个 `delta.content`
> 和 `delta.tool_calls`。转换器需**自行合成**:在第一次出现 `delta.content` 时发出
> `ContentBlockStart(0, Text)`,在 `delta.tool_calls` 出现时发出
> `ContentBlockStart(n, ToolUse { id, name })`。注意 content 和 tool_calls 在同一 chunk **互斥**
> (需验证),合成规则为"来什么类型,合什么边界"。
> 3. **Tool call 全局序号映射**OpenAI 的 `delta.tool_calls[i].index` 是**该 chunk 内的局部索引**,
> 需要累积状态(`HashMap<局部index, 全局index>`)来推导全局 content block 序号。
> 对比:**已推演**的 IR 设计将 ToolCallStart 合并到 ContentBlockStartindex 在 IR 层面统一为
> 全局序号,差异完全封装在 Provider 层。
> 4. **多 Choice 的处理**:当前代码只处理 `choices[0]`。新设计中是继续忽略其他 choice 还是
> 通过某种机制暴露(如 `extra` 中携带)?
> 5. **usage 的时机**OpenAI 的 usage 通常在最后一个 chunk 中携带,与 finish_reason 在同一 chunk。
> 是先发 `CostUpdate(PartialUsage{...})` 再发 `MessageComplete`,还是反过来?
> **需要推演:**
> - `OpenaiStreamToEvents<S>` 转换器的状态机设计(跟踪的局部状态、事件产出规则)
> - 边界情况:SSE 行乱序、`[DONE]` 标记的处理、网络断开重连
> **✅ 推演结论(2026-06-18):**
>
> ### 架构分层
>
> OpenAI 的流式转换拆为**两层**,字节解析层通用、事件转换层 OpenAI 特定:
>
> ```
> bytes_stream()
>
>
> SseByteStream [sse.rs — 通用层]
> │ 逐行分割、去 "data:" 前缀、过滤 "[DONE]"、缓冲拼接碎片
>
> Stream<Item=Result<String, LlmError>> ← JSON 字符串行
>
>
> OpenaiStreamToEvents [openai/stream.rs — OpenAI 特定]
> │ 反序列化为 OpenaiChatChunk → 状态机转换
>
> Stream<Item=Result<StreamEvent, LlmError>> ← IR 语义事件
> ```
>
> **Layer 1 — `SseByteStream<S>`(新增 `src/llm/provider/sse.rs`**
>
> 从当前 `SseChunkStream``openai.rs` 内联实现)提取字节解析逻辑为通用 SSE 解析器。
> 产出 JSON 字符串行,不绑定 OpenAI 格式。Anthropic 直接复用。
>
> **Layer 2 — `OpenaiStreamToEvents`(新增 `src/llm/provider/openai/stream.rs`**
>
> 接收 JSON 字符串行,反序列化为 `OpenaiChatChunk`(保留作为内部格式),
> 通过状态机转换为 `StreamEvent` 事件流。
>
> ### 状态机设计
>
> ```rust
> pub struct OpenaiStreamToEvents<S> {
> inner: S, // Stream<Item=Result<String, LlmError>>
> // ── 状态 ──
> global_index: u32, // 下一个可用 content block 序号
> current_block: Option<(u32, CurrentBlockType)>, // 当前活跃 block
> tool_call_indices: HashMap<u32, u32>, // OpenAI tool_call.index → 全局序号
> is_complete: bool,
> }
>
> enum CurrentBlockType { Text, Refusal, ToolUse }
> ```
>
> **每条 JSON 行的处理流程:**
>
> ```
> 收到一行 JSON 字符串
> ├─ 解析为 OpenaiChatChunk
>
> ├─ Phase 1: 处理 delta 内容(先增量)
> │ ├─ delta.content → ensure_block(Text) → TextDelta
> │ ├─ delta.refusal → ensure_block(Refusal) → RefusalDelta
> │ └─ delta.tool_calls
> │ ├─ 新 tool callfunction.name 有值)
> │ │ → ensure_block(ToolUse, id, name) → (不产参数事件,等后续 arguments)
> │ └─ 已有 tool callfunction.arguments 有值)
> │ → ToolCallArgumentsDelta(index, arguments)
>
> └─ Phase 2: 处理汇总(后收束)
> ├─ finish_reason 存在 → close_current_block() + MessageComplete
> ├─ usage 存在 → CostUpdate
> └─ 同时存在 → CostUpdate → close_block → MessageComplete
> ```
>
> **核心抽象 `ensure_block`** 当新 chunk 的 delta 类型与当前活跃 block 不同时,
> 自动关闭当前 blockemit `ContentBlockEnd`)并开启新 blockemit `ContentBlockStart`)。
> 同类型继续时只发增量事件,不切换。
>
> **状态转移表:**
>
> | 当前状态 | 收到 delta.content | 收到 delta.tool_calls (new) | 收到 delta.tool_calls (延续) | 收到 finish_reason |
> |----------|-------------------|----------------------------|----------------------------|-------------------|
> | 无活跃 block | → ContentBlockStart(Text)<br>→ TextDelta | → ContentBlockStart(ToolUse) | (不应发生) | → MessageComplete |
> | Text 活跃中 | → TextDelta | → ContentBlockEnd<br>→ ContentBlockStart(ToolUse) | (不应发生,与 content 互斥) | → ContentBlockEnd<br>→ CostUpdate<br>→ MessageComplete |
> | ToolUse 活跃中 | → ContentBlockEnd<br>→ ContentBlockStart(Text) | → ContentBlockEnd<br>→ ContentBlockStart(new ToolUse) | → ToolCallArgumentsDelta | → ContentBlockEnd<br>→ CostUpdate<br>→ MessageComplete |
>
> ### 各议题结论
>
> | # | 议题 | 结论 | 理由 |
> |---|------|------|------|
> | 1 | **SSE 字节解析复用** | 提取为通用 `SseByteStream``sse.rs`),`OpenaiStreamToEvents` 在其上封装 | Anthropic 复用相同字节协议、分层测试、职责清晰 |
> | 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,非目标已明确 |
> | 5 | **Usage 时机** | 先发 `CostUpdate` → 再发 `MessageComplete`,同一 chunk 内串行 | 汇聚算法兼容两者顺序,但语义上先用量后完成更合理 |
> | 6 | **ToolCallEnd 发出** | OpenAI 层不显式发出 `ToolCallEnd`,依赖汇聚算法 `finalize()` 兜底 | `ToolCallEnd` 保留给 Anthropic`content_block_stop` 场景);`ContentBlockEnd` 已标 block 完成,`finalize` 时 `tool_call_args` 已累积完整 |
> | 7 | **Refusal 处理** | 检测 `delta.refusal`emit `ContentBlockStart(Refusal)` + `RefusalDelta` + `ContentBlockEnd` | OpenAI 特有字段,IR 已有 `ContentBlockType::Refusal` |
> | 8 | **代码消重** | 删除 `stream.rs::parse_chunk_stream`(无人调用);`cycle.rs::submit_stream` 直接消费 `Result<StreamEvent>` 流 | 转换逻辑统一到 Provider 层,LlmCycle 只负责编排和 hook |
>
> ### 边界情况
>
> | 场景 | 处理方式 |
> |------|---------|
> | **`[DONE]` 行** | `SseByteStream` 层过滤,不传递到事件层(当前已有逻辑) |
> | **delta 为空 + finish_reason** | 只处理 Phase 2,关闭当前 block 后 emit MessageComplete |
> | **同一 chunk 含 delta.content + finish_reason** | Phase 1 先处理 delta 发 TextDeltaPhase 2 关闭 block 发 Complete |
> | **同一 chunk 含 delta.tool_calls + finish_reason** | 先处理所有 tool calls(映射 + arguments 累积),再关 block |
> | **usage 单独 chunk 下发** | 触发 Phase 2 但 finish_reason 为 None → 只发 CostUpdate,不发 MessageComplete |
> | **网络断开** | reqwest 返回 Err → emit StreamEvent::Erroris_complete = true |
> | **JSON 解析失败** | emit StreamEvent::Error,终止流(防御性处理) |
>
> ### 文件变更清单
>
> | 操作 | 文件 | 说明 |
> |------|------|------|
> | 新增 | `src/llm/provider/sse.rs` | 通用 SSE 字节解析层,从当前 `openai.rs` 的 `SseChunkStream` 提取核心逻辑 |
> | 新增 | `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>` |
> | 修改 | `src/llm/cycle.rs` | `submit_stream` 直接消费 `StreamEvent` 流,移除内联 chunk→event 转换 |
> | 修改 | `src/llm/stream.rs` | 删除 `parse_chunk_stream`(无人调用) |
>
> 优先级:高(Phase 2 OpenAI Provider 重构的核心任务)
### 5.2 AnthropicProviderMessages API