From b7d9a41f94240f5267cd4750aa40e3a1b0f3ffaa Mon Sep 17 00:00:00 2001 From: xiaoxixi Date: Fri, 17 Jul 2026 12:57:35 +0800 Subject: [PATCH] docs: define streaming turn architecture --- docs/ARCHITECTURE.md | 2 + docs/STREAMING_TURN_DESIGN.md | 807 ++++++++++++++++++++++++++++++++++ 2 files changed, 809 insertions(+) create mode 100644 docs/STREAMING_TURN_DESIGN.md diff --git a/docs/ARCHITECTURE.md b/docs/ARCHITECTURE.md index ac2dfe6..595fbb8 100644 --- a/docs/ARCHITECTURE.md +++ b/docs/ARCHITECTURE.md @@ -2,6 +2,8 @@ 本文档描述 PicoBot 当前实现的运行时边界、数据流、并发模型和演进约束。它面向维护者和后续参与改进的 Agent,是代码架构的主入口;行为细节仍以代码和测试为最终依据。 +拟议中的流式模型输出、reasoning 展示、活动 Turn 快照和 Channel 实时投递架构见 [STREAMING_TURN_DESIGN.md](STREAMING_TURN_DESIGN.md)。该文档是尚未实现的设计,不代表当前运行时行为。 + ## 1. 设计目标 PicoBot 是一个单进程、异步、可扩展的个人 AI 助手运行时。核心目标是: diff --git a/docs/STREAMING_TURN_DESIGN.md b/docs/STREAMING_TURN_DESIGN.md new file mode 100644 index 0000000..91c2159 --- /dev/null +++ b/docs/STREAMING_TURN_DESIGN.md @@ -0,0 +1,807 @@ +# 流式 Turn、Reasoning 展示与 Channel 投递设计 + +> 状态:拟议设计,尚未实现。 +> +> 本文定义 PicoBot 未来的流式模型输出、reasoning 展示、工具过程展示和 Channel 实时投递架构。当前运行时行为仍以 `docs/ARCHITECTURE.md` 和代码为准。实施本设计时允许重做内部、WebSocket、TUI、WebUI 和 Channel 协议;只要求 SQLite 数据库提供可靠迁移路径。 + +## 1. 背景 + +当前 Provider 只提供一次性 `chat()` 调用。AgentLoop 等待完整响应,将 `content`、`reasoning_content` 和 tool calls 组装为 `ChatMessage`,Session 在 AgentLoop 完成后原子持久化本轮消息,再通过 `OutboundMessage` 发送最终正文。 + +这个模型具有清晰的持久化语义,但无法表达: + +- 正文和 reasoning 的实时增量; +- reasoning、正文、工具调用在一个 Turn 中的自然交错; +- TUI 和 WebUI 对同一个运行中 Turn 的一致展示; +- 飞书等 Channel 通过编辑消息呈现流式效果; +- 不支持编辑的 Channel 自动降级为只发送最终结果; +- 取消、失败、慢消费者和投递失败时的确定行为。 + +现有 `Channel::send_delta(chat_id, delta)` 没有 Turn 身份、消息身份、reasoning/text 分类、工具边界、终态和取消语义,也绕过现有出站排序机制,不适合作为后续架构基础。 + +## 2. 参考实现结论 + +本设计综合了 `reference/ryvos`、`reference/zeroclaw`、`reference/hermes-agent` 和 PicoBot 当前实现。 + +### 2.1 Ryvos + +Ryvos 把 thinking 作为正式内容块,并区分 `TextDelta` 与 `ThinkingDelta`。它说明 Provider 层必须结构化解析 reasoning、正文和工具调用,不能依赖最终字符串中的 `` 后处理。 + +可吸收: + +- Provider 流的类型化增量; +- reasoning 与正文分离; +- 工具参数增量组装; +- thinking-only 响应的明确处理; +- reasoning effort/thinking budget 的统一配置概念。 + +### 2.2 ZeroClaw + +ZeroClaw 把 reasoning 作为不透明 Provider 数据保留,用于要求历史回放的模型,同时处理 `reasoning_content`、`reasoning`、内联 `` 和不同 Provider 的回放限制。 + +可吸收: + +- 可展示 reasoning 与 Provider 回放状态分离; +- reasoning 字段别名归一化; +- Provider 专用历史状态不能跨 Provider 发送; +- 内联 think block 必须在 Provider 归一化边界处理; +- reasoning 是否展示与是否回放是两个独立策略。 + +### 2.3 Hermes + +Hermes 新增了 Agent 到 Gateway 的结构化展示事件,并明确规定流事件属于 presentation,而不是 conversation history。它还通过 message segment boundary 处理“工具前正文 → 工具 → 工具后正文”。 + +可吸收: + +- 流事件描述发生的事实,不携带平台发送策略; +- 展示流与持久化历史严格分离; +- 工具边界必须结束当前正文 segment; +- 高频更新需要合并,关键事件发送前需要 flush; +- TUI 的活动 Turn 与已完成 transcript 分离; +- Channel/平台决定如何呈现统一状态。 + +不直接复制: + +- typed events、旧 callbacks 和 TUI 字符串事件并存; +- reasoning 使用独立 callback,没有进入新事件模型; +- 客户端用大型 TurnController 重建服务端状态; +- 一个 GatewayStreamConsumer 同时承担聚合、限流、平台编辑、think 清理、overflow 和 fallback。 + +### 2.4 PicoBot + +PicoBot 已有以下适合保留的不变量: + +- 同一 Session 由单 worker 串行处理; +- 不同 Session 并发; +- `worker_generation` 和 `state_version` 防止迟到结果提交; +- 完整 Agent Turn 通过原子持久化接口提交; +- 普通出站消息按 `(channel, chat_id)` 有序投递; +- Channel、Session、Agent、Provider 和 Storage 边界明确。 + +流式设计不能破坏这些不变量。 + +## 3. 设计目标 + +1. OpenAI-compatible 和 Anthropic Provider 支持流式正文、reasoning、tool calls 和 usage。 +2. TUI 与 WebUI 使用同一运行态模型显示正文、reasoning、工具进度和取消/失败状态。 +3. Channel 可以选择实时更新或只发送最终结果。 +4. 支持消息编辑的 Channel 能以同一远端消息呈现流式效果。 +5. 慢客户端或慢 Channel 不得反压 Provider 和 AgentLoop,也不得积压大量过时 token。 +6. 丢失任意中间更新后,下一次更新必须自动收敛到正确状态。 +7. `Completed` 必须表示本轮数据库提交已经成功。 +8. reasoning 展示文本与 Provider 回放状态必须隔离。 +9. Channel 展示差异不能改变模型上下文和数据库历史。 +10. 保持模块数量、事件词汇和状态 owner 尽可能少。 + +## 4. 非目标 + +- 不保留旧 WebSocket、TUI、WebUI 或 Channel 流式协议兼容性。 +- 不要求每个 Provider 都能返回可展示 reasoning。 +- 不把 Provider 的加密/签名 reasoning payload 展示给用户。 +- 不逐 token 持久化数据库。 +- 不保证重连后恢复尚未完成 Turn 的每一个历史帧。 +- 不让所有外部 Channel 默认以多条追加消息模拟流式效果。 +- 不把流式展示事件作为可重放的事件溯源日志。 + +## 5. 核心决策 + +### 5.1 只有两个权威模型 + +系统只维护两个跨层权威模型: + +- `ConversationMessage`:最终持久化事实,用于会话历史和下轮模型上下文; +- `TurnState`:单次运行的临时展示状态。 + +Provider 的 SSE chunk 是 Provider 内部输入;Channel 的远端消息 ID 是单个 TurnSink 的私有投递状态。两者都不是全局领域模型。 + +### 5.2 服务端拥有唯一运行态 + +Session 侧 `TurnController` 是运行中 Turn 的唯一状态 owner。TUI、WebUI 和 Channel 不根据一串增量自行重建 reasoning、正文、工具和 segment 关系,只渲染服务端发布的 `TurnSnapshot`。 + +### 5.3 下发幂等快照,不下发可靠 token 流 + +Provider 到 AgentLoop 使用 delta;TurnController 到展示端使用包含完整当前状态的快照。 + +每个快照带单调递增 `revision`。消费者只接受 revision 更大的快照,并以新快照整体替换旧状态。因此: + +- 丢失中间更新不会损坏内容; +- 慢消费者可以跳过过时状态; +- 更新重试不会重复拼接正文; +- 最终快照能修复暂态渲染; +- 外部 Channel 编辑消息天然获得完整累计内容。 + +### 5.4 latest-wins,而不是 token 队列 + +每个活动 Turn 使用 Tokio `watch` 或等价 latest-value primitive 发布 `Arc`。生产者覆盖旧值,消费者读取最新值。终态显式编码在快照中,不能仅依赖 sender 关闭表达完成。 + +### 5.5 每个 Channel Turn 使用独立 TurnSink + +Channel 为每次 Turn 创建一个 sink。Sink 独占远端消息 ID、编辑状态和平台私有资源,终态后销毁。通用协调器不保存平台消息映射,也不理解飞书卡片 API。 + +### 5.6 完整消息投递与活动 Turn 展示分离 + +- `MessageBus` / `OutboundDispatcher`:完整消息、通知、命令结果和需要可靠确认的独立投递; +- `DeliveryCoordinator`:活动 Turn 快照、展示策略、节流和 TurnSink 生命周期。 + +不把 token 或快照塞入现有 outbound MPSC。 + +## 6. 数据模型 + +### 6.1 Provider 私有流 + +```rust +pub enum ProviderChunk { + Text(String), + Reasoning(String), + ToolCallStart { + index: usize, + id: Option, + name: Option, + }, + ToolCallArguments { + index: usize, + delta: String, + }, + ProviderState(ProviderReasoningState), + Usage(Usage), + Done(FinishReason), +} +``` + +`ProviderChunk` 只允许在 `providers` 与 `agent` 模块间使用,不进入 Bus、Session 协议或 Channel。 + +### 6.2 可展示 reasoning 与回放状态 + +```rust +pub struct ProviderReasoningState { + pub provider: String, + pub payload: serde_json::Value, +} + +pub struct AssistantMessageData { + pub content: String, + pub reasoning: Option, + pub provider_state: Option, + pub tool_calls: Vec, +} +``` + +约束: + +- `reasoning` 可以按展示策略下发; +- `provider_state` 永远不下发给客户端或 Channel; +- `provider_state.provider` 与当前 Provider 不一致时禁止回放; +- 压缩历史时默认不把原始 reasoning 写入 Timeline; +- 日志不得记录完整 reasoning 或 provider payload。 + +### 6.3 Turn 标识 + +```rust +pub struct TurnId(pub uuid::Uuid); +pub struct BlockId(pub uuid::Uuid); +``` + +一次用户输入对应一个 Turn。Turn 开始时预分配最终 assistant `message_id`,使运行态和最终历史能稳定关联。 + +### 6.4 TurnState + +```rust +pub struct TurnState { + pub id: TurnId, + pub session_id: String, + pub message_id: String, + pub revision: u64, + pub status: TurnStatus, + pub phase: TurnPhase, + pub blocks: Vec, + pub usage: Option, + pub error: Option, +} + +pub enum TurnStatus { + Running, + Completed, + Cancelled, + Failed, +} + +pub enum TurnPhase { + Queued, + Reasoning, + Responding, + Acting, + Finalizing, +} +``` + +`TurnPhase` 是展示状态,不等于模型 reasoning: + +- 等待首个模型 chunk 时可显示 `Queued`; +- 收到 reasoning delta 时进入 `Reasoning`; +- 收到正文 delta 时进入 `Responding`; +- 执行工具时进入 `Acting`; +- 模型结束、等待持久化时进入 `Finalizing`。 + +### 6.5 有序 TurnBlock + +```rust +pub enum TurnBlock { + Reasoning { + id: BlockId, + iteration: u32, + text: String, + }, + Assistant { + id: BlockId, + iteration: u32, + text: String, + }, + Tool { + id: String, + iteration: u32, + name: String, + arguments: serde_json::Value, + status: ToolStatus, + preview: Option, + }, +} + +pub enum ToolStatus { + Running, + Completed, + Failed, +} +``` + +有序 block 直接表达 reasoning、正文和工具的交错,不再维护多个平行字符串或让客户端猜测工具边界。 + +### 6.6 Agent 语义事件 + +```rust +pub enum TurnEvent { + ReasoningDelta { + iteration: u32, + delta: String, + }, + TextDelta { + iteration: u32, + delta: String, + }, + TextSegmentFinished { + iteration: u32, + }, + ToolStarted { + iteration: u32, + call: ToolCall, + }, + ToolFinished { + iteration: u32, + call_id: String, + success: bool, + preview: Option, + }, +} +``` + +AgentLoop 只发过程事实。Turn 的 start、finalize、complete、cancel 和 fail 由 Session worker 调用 TurnController,因为 Session 才拥有生命周期、持久化和 stale-state 判断。 + +## 7. Provider 层 + +### 7.1 流式优先接口 + +```rust +#[async_trait] +pub trait LLMProvider: Send + Sync { + async fn stream( + &self, + request: ChatCompletionRequest, + ) -> Result; +} +``` + +标题生成、压缩等需要完整响应的代码通过 `collect_provider_stream()` 收集同一实现,避免分别维护 stream 和 non-stream HTTP 路径。 + +### 7.2 OpenAI-compatible + +至少处理: + +- `delta.content`; +- `delta.reasoning_content`; +- `delta.reasoning`; +- 同一 payload 同时出现 content 与 reasoning; +- tool call id/name/arguments 分片; +- usage-only final chunk; +- reasoning-only 响应; +- 内联 think tag 被任意 SSE chunk 切分。 + +``、`` 等内联标签使用有状态 parser 在 Provider 归一化边界转换为 `ProviderChunk::Reasoning`,不能在 Channel 或客户端重复清理。 + +### 7.3 Anthropic + +至少处理: + +- text content block; +- thinking/redacted thinking block; +- signature 或其他回放元数据; +- tool_use block 和 input JSON delta; +- content block start/delta/stop; +- message usage 和 stop reason。 + +thinking 文本进入 `reasoning`,签名和原始 block 进入 `provider_state`。历史回放必须保持 Provider 要求的块顺序和签名完整性。 + +### 7.4 reasoning-only + +Provider 不擅自把 reasoning 提升为正文。AgentLoop 在最终轮发现只有 reasoning、没有正文且没有工具调用时,执行统一产品策略: + +- 默认生成一个明确的空正文终态,并在 UI 保留 reasoning;或 +- 后续配置允许把 reasoning 作为正文 fallback。 + +首版建议不泄漏 reasoning 为正文,由 UI 显示“模型未返回最终正文”。 + +## 8. AgentLoop 与 TurnController + +### 8.1 AgentLoop + +AgentLoop 消费 `ProviderChunk` 并: + +- 累积本轮完整 AssistantMessageData; +- 将展示事实发给 `TurnEmitter`; +- 组装 tool calls; +- Provider 完成当前迭代后执行工具; +- 在工具开始前发出 `TextSegmentFinished`; +- 将完整 assistant/tool messages 加入内部历史; +- 返回最终 `AgentProcessResult`。 + +AgentLoop 不访问 MessageBus、DeliveryCoordinator 或 Channel。 + +### 8.2 TurnController + +TurnController: + +- 是 TurnState 的唯一写入者; +- 把 TurnEvent reduce 为有序 blocks; +- 维护 revision、status 和 phase; +- 发布最新 TurnSnapshot; +- 不执行 Provider、工具、数据库或 Channel I/O。 + +推荐接口: + +```rust +impl TurnController { + pub fn start(... ) -> (Self, TurnEmitter, watch::Receiver>); + pub fn begin_finalizing(&mut self); + pub fn complete(&mut self, usage: Option); + pub fn cancel(&mut self, reason: Option); + pub fn fail(&mut self, error: String); +} +``` + +`TurnEmitter` 应轻量、无 Channel 感知,并在 Turn 终止后拒绝新事件。 + +### 8.3 stale-state + +Session worker 在以下位置校验 `worker_generation` 和必要的 `state_version`: + +- 创建 Turn 后、调用 Provider 前; +- 发布会产生用户可见变化的快照前; +- 工具批次完成后; +- 最终数据库提交前; +- 发布 Completed 前。 + +旧 generation 的 Turn 必须进入 Cancelled 或静默终止,禁止继续编辑 Channel 远端消息。 + +## 9. DeliveryCoordinator + +DeliveryCoordinator 是活动 Turn 的唯一展示协调器,职责包括: + +1. 订阅 `watch::Receiver`; +2. 解析当前目标的 PresentationPolicy; +3. 在下发前移除隐藏的 reasoning/tool blocks; +4. 根据 Channel LivePolicy 节流; +5. 为 Turn 创建并持有 TurnSink; +6. 对 Running 快照进行 best-effort 更新; +7. 对终态快照立即 flush; +8. 有界等待 sink 结束; +9. 报告最终投递结果。 + +### 9.1 合并和背压 + +- `watch` 自动覆盖过时快照; +- WebSocket/cli_chat 建议最多约 30 FPS; +- 飞书建议从 500ms 更新间隔开始; +- 正在滚动或渲染压力高时,客户端无需向服务端反馈节流,丢弃中间快照即可; +- 终态不等待节流 timer,必须立即发送; +- Running 更新失败不重试旧快照,等待下一最新快照; +- Completed 最终投递使用正常可靠重试语义。 + +### 9.2 与 OutboundDispatcher 的关系 + +DeliveryCoordinator 不取代 OutboundDispatcher。 + +| 组件 | 负责 | +|------|------| +| OutboundDispatcher | 完整独立消息、通知、命令结果、定时投递、可靠重试 | +| DeliveryCoordinator | 活动 Turn 的运行快照、节流、展示过滤和 sink 生命周期 | + +两者对同一 `(channel, chat_id)` 的最终写操作必须有统一排序边界。实现时可以复用 per-conversation lane owner,但不能把所有快照排进 lane 的普通 MPSC;lane 应只持有 TurnSink 任务或 latest snapshot receiver。 + +## 10. Channel 与 TurnSink + +### 10.1 接口 + +删除 `Channel::send_delta`,保留普通 `send`,增加: + +```rust +pub enum LivePolicy { + FinalOnly, + Snapshot { + min_interval: Duration, + }, +} + +#[async_trait] +pub trait Channel: Send + Sync + 'static { + fn live_policy(&self) -> LivePolicy; + + async fn open_turn( + &self, + target: TurnTarget, + ) -> Result, ChannelError>; + + async fn send(&self, msg: OutboundMessage) -> Result<(), ChannelError>; +} + +#[async_trait] +pub trait TurnSink: Send { + async fn update(&mut self, snapshot: &TurnSnapshot) -> Result<(), ChannelError>; + async fn finish(self: Box, snapshot: &TurnSnapshot) -> Result<(), ChannelError>; + async fn abort(self: Box, snapshot: &TurnSnapshot) -> Result<(), ChannelError>; +} +``` + +`TurnSink` 的每个调用都接收完整、过滤后的快照。Sink 不拼接 token。 + +### 10.2 cli_chat + +- `LivePolicy::Snapshot`,默认约 33ms; +- update/finish/abort 都发送统一 `turn_updated` frame; +- WebUI 与 TUI 使用完全相同的 TurnSnapshot; +- Session 不匹配时客户端忽略渲染,但可标记未读; +- 终态后客户端可请求历史作最终校准。 + +### 10.3 飞书 + +- 配置关闭实时展示时使用 FinalOnly sink; +- 开启时使用 Snapshot sink; +- 第一个有可见内容的 Running 快照创建卡片; +- 后续快照编辑同一卡片; +- sink 内持有远端 message ID; +- reasoning/tool block 由 PresentationPolicy 决定是否进入卡片; +- 中间编辑失败不影响 Agent; +- finish 做最终编辑,必要时退化为发送一条完整最终消息; +- 卡片长度限制、拆分和平台限流属于 FeishuTurnSink 私有实现。 + +### 10.4 不支持编辑的 Channel + +实现 FinalOnlyTurnSink:忽略 Running 快照,只在 finish 时发送最终投影。默认不通过多条追加消息模拟流式,以免产生无法收回的碎片消息。 + +## 11. 展示策略 + +```rust +pub struct PresentationPolicy { + pub live: bool, + pub reasoning: ReasoningVisibility, + pub tools: ToolVisibility, +} + +pub enum ReasoningVisibility { + Hidden, + Collapsed, + Expanded, +} + +pub enum ToolVisibility { + Hidden, + Compact, + Detailed, +} +``` + +建议默认值: + +- TUI/WebUI:live=true,reasoning=Collapsed,tools=Detailed; +- 外部 Channel:live 由 Channel 配置决定,reasoning=Hidden,tools=Compact; +- Scheduler/无人值守投递:FinalOnly,reasoning=Hidden。 + +策略由 DeliveryCoordinator 在数据离开 Gateway 核心前应用。Channel 和客户端不能只靠“隐藏 UI”实现 reasoning 保密。 + +## 12. TUI 与 WebUI + +客户端状态简化为: + +```text +history 已持久化消息 +active_turn 当前 TurnSnapshot(每个 session 最多一个) +``` + +收到快照时: + +```text +if snapshot.revision > active_turn.revision: + active_turn = snapshot +``` + +渲染规则: + +- Reasoning block 显示为折叠或展开区域; +- Assistant block 显示为正文 segment; +- Tool block 显示运行中/成功/失败状态; +- phase 控制 spinner 文案; +- Completed/Cancelled/Failed 显示明确终态; +- 当前 session 之外的快照不进入当前消息列表; +- Completed 后以历史响应替换 active turn,避免展示态与数据库态长期并存。 + +WebUI 对 Running Markdown 可以按动画帧或快照频率渲染,Completed 时进行最终 sanitize。TUI 原地重绘 active turn,不把每个更新追加为新 transcript 行,也不在用户向上滚动时强制跳到底部。 + +## 13. WebSocket 协议 + +不保留旧 `assistant_response` 流程。活动 Turn 使用一个统一 frame: + +```json +{ + "type": "turn_updated", + "turn": { + "id": "...", + "session_id": "...", + "message_id": "...", + "revision": 12, + "status": "running", + "phase": "responding", + "blocks": [] + } +} +``` + +Completed、Cancelled 和 Failed 仍使用 `turn_updated`,只改变完整快照的 status。避免为每个生命周期阶段增加一组容易漂移的 frame 类型。 + +历史协议应返回持久化后的 reasoning、completion_status、turn_id 和 iteration,但不返回 provider_state。 + +## 14. 持久化和提交顺序 + +流式过程中不逐 token 写 SQLite。完成顺序固定为: + +```text +Provider 完成 + → AgentLoop 组装 emitted_messages + → Session 校验 generation/state_version + → 原子写入本轮全部消息 + → TurnController.complete() + → 发布 Completed 快照 + → TurnSink.finish() +``` + +因此 `TurnStatus::Completed` 的含义是:数据库已经提交成功,最终展示可以安全收敛到历史。 + +### 14.1 取消 + +建议首版语义: + +- 没有 Assistant 正文:不持久化 assistant 消息,Turn 标记 Cancelled; +- 已向用户展示部分正文:持久化部分正文并标记 `completion_status=cancelled`; +- 已完成的 assistant/tool/tool-result 链必须保持 Provider 可接受的结构; +- reasoning 可随取消消息保存,但展示仍受 policy 控制; +- 取消后 TurnEmitter 关闭,迟到 delta 被丢弃。 + +### 14.2 失败 + +- Provider 在任何可见正文前失败:Turn Failed,错误作为结构化 error 展示,不创建 assistant 历史; +- 已产生部分正文后失败:按 interrupted partial 保存,标记 `completion_status=interrupted`; +- 持久化失败:不得发布 Completed,Turn Failed,并明确告知用户流式预览未保存; +- Running Channel 更新失败不改变 Turn 结果;最终 finish 失败按现有投递错误处理。 + +## 15. SQLite 迁移 + +建议为 `messages` 增加: + +```text +turn_id TEXT NULL +iteration INTEGER NULL +completion_status TEXT NOT NULL DEFAULT 'completed' +reasoning_content TEXT NULL -- 已存在,语义调整为可展示 reasoning +provider_state TEXT NULL -- JSON,Provider 私有回放状态 +``` + +要求: + +- 新库 schema 测试; +- 旧库迁移测试; +- 已有 `reasoning_content` 数据原样保留; +- `provider_state` 解析失败时降级为不回放,不能使历史不可读; +- Session 加载继续修复 tool-call chains; +- 原子提交覆盖完整 Turn 的所有 emitted messages。 + +是否新增 `turns` 表首版暂不需要。运行中 Turn 只存在内存,历史可通过 messages.turn_id 分组。如果未来要跨 Gateway 重启恢复运行态,再单独设计 durable turn lease/state。 + +## 16. 生命周期与并发不变量 + +1. 每个 Session 最多一个活动主 Turn。 +2. TurnController 是 TurnState 的唯一写入者。 +3. AgentLoop、DeliveryCoordinator 和 TurnSink 不持有 Session mutex 执行慢 I/O。 +4. `worker_generation` 变化后,旧 Turn 不得发布新快照或提交消息。 +5. Running 快照是 best-effort;终态快照必须显式、完整且有界投递。 +6. 慢 sink 只能跳过中间状态,不能阻塞 Provider 或 AgentLoop。 +7. Completed 必须晚于数据库成功提交。 +8. TurnSink 的生命周期由 DeliveryCoordinator 所有,并通过 TaskSupervisor 回收。 +9. Provider stream、节流 timer、Channel 编辑、取消和 shutdown 都必须有硬时间界限。 +10. presentation 过滤不能修改 ConversationMessage 或 Agent 历史。 + +## 17. 模块布局 + +```text +src/providers/stream.rs + ProviderChunk、ProviderStream、FinishReason、collect helper + +src/agent/turn_event.rs + TurnEvent、TurnEmitter + +src/session/turn.rs + TurnState、TurnBlock、TurnController、TurnSnapshot + +src/delivery/mod.rs +src/delivery/coordinator.rs +src/delivery/policy.rs + watch 订阅、展示过滤、节流、sink 生命周期 + +src/channels/base.rs + Channel、LivePolicy、TurnSink + +src/channels/cli_chat.rs + WebSocketTurnSink + +src/channels/feishu.rs + FeishuTurnSink + +src/protocol.rs + TurnSnapshot 序列化 +``` + +`observability` 继续记录 Agent/tool 遥测,不承担 UI stream。`MessageBus` 不新增 token/turn 队列。 + +## 18. 明确拒绝的替代方案 + +### 18.1 每 token 一个 OutboundMessage + +拒绝原因:填满 bounded bus/lane、重试乱序、慢 Channel 反压 Agent、最终消息和暂态更新语义混淆。 + +### 18.2 端到端 delta 协议 + +拒绝原因:客户端和每个 Channel 都必须实现累积、去重、segment、取消和丢帧恢复状态机,最终产生多个事实 owner。 + +### 18.3 客户端自行组合 reasoning/tool/text + +拒绝原因:TUI、WebUI 和 Channel 行为会漂移;重连和切 session 时难以恢复;服务端已经拥有全部事实。 + +### 18.4 把流式事件写入数据库 + +拒绝原因:消息历史膨胀,事务语义复杂,压缩和 Provider 回放被展示细节污染。 + +### 18.5 在 Channel 全局保存 turn_id 映射 + +拒绝原因:owner 和清理边界不清晰。每 Turn 一个 sink 可以让远端消息状态自然随生命周期释放。 + +### 18.6 一个巨型跨平台 StreamConsumer + +拒绝原因:通用合并/策略与平台 API 细节耦合。DeliveryCoordinator 只做统一调度,具体远端编辑由各 TurnSink 自己实现。 + +## 19. 验证策略 + +### 19.1 Provider + +- SSE 任意字节和 UTF-8 边界切分; +- reasoning/content 同时出现; +- reasoning 与 content 字段别名; +- think tag 跨 chunk; +- 多 tool calls 交错参数 delta; +- usage-only chunk; +- Anthropic thinking signature round-trip; +- 中途断线、超时和取消。 + +### 19.2 TurnController + +- reasoning → text → tool → reasoning → text 的 block 顺序; +- segment boundary; +- parallel tool status; +- revision 严格递增; +- 终态后拒绝新事件; +- reasoning-only、empty response、取消和失败。 + +建议用 property tests 验证:任意合法 TurnEvent 序列 reduce 后不产生相邻可合并同类 block、重复 tool id 或终态后变更。 + +### 19.3 DeliveryCoordinator + +- 慢 sink 只收到最新快照; +- 终态绕过节流; +- hidden reasoning 在到达 sink 前已经移除; +- Running 更新失败后能由下一快照恢复; +- finish 可靠重试; +- shutdown 有界; +- stale generation 停止更新。 + +### 19.4 客户端和 Channel + +- WebUI/TUI revision 去重; +- session 切换不显示迟到 Turn; +- Completed 后历史校准; +- Markdown 未完成块与最终块; +- 飞书 create/edit/final fallback; +- FinalOnly sink 不发送中间内容; +- 远端长度限制与节流。 + +### 19.5 Storage + +- 新库 schema; +- 旧 schema 迁移; +- provider_state 损坏降级; +- cancelled/interrupted 消息恢复; +- 完整 Turn 原子提交失败不产生部分历史。 + +## 20. 实施顺序 + +1. 增加数据库迁移和新的消息 reasoning/provider state 语义。 +2. 引入 ProviderChunk 和流式优先 Provider trait。 +3. 实现 OpenAI-compatible streaming 和 collect helper。 +4. 实现 TurnEvent、TurnController、TurnSnapshot 及纯单元测试。 +5. AgentLoop 接入 TurnEmitter,但暂时只使用最终结果投递。 +6. 增加 DeliveryCoordinator、LivePolicy 和 TurnSink。 +7. 重做 cli_chat WebSocket 为统一 `turn_updated` 快照。 +8. 重做 TUI/WebUI 为 `history + active_turn` 渲染模型。 +9. 实现 Anthropic streaming 和完整 provider_state 回放。 +10. 实现 FeishuTurnSink 的创建、编辑、终态和 fallback。 +11. 删除旧 `send_delta`、旧 reasoning 特例和 Channel 侧 think 清理。 +12. 更新 `ARCHITECTURE.md`、README、AGENTS.md 和运行时产品知识。 + +每一步都必须保持非流式最终回复可用;但不为旧协议保留双栈兼容代码。 + +## 21. 实施前需确认的产品策略 + +这些选择不改变架构,但必须在实现前确定默认值: + +1. `/stop` 后是否持久化已展示的部分正文。本文建议持久化并标记 cancelled。 +2. TUI/WebUI reasoning 默认折叠还是隐藏。本文建议折叠。 +3. 外部 Channel reasoning 默认策略。本文建议隐藏。 +4. reasoning-only 是否允许提升为正文。本文建议不提升。 +5. provider_state 的保留期限和 LLM 原始响应审计策略。 + +## 22. 架构验收标准 + +设计完成实现后,应能用以下陈述准确描述系统: + +- Provider 只负责模型协议,AgentLoop 只负责模型/工具语义。 +- Session 拥有 Turn 生命周期和最终持久化。 +- TurnController 是运行态的唯一事实来源。 +- DeliveryCoordinator 只投影展示,不修改历史。 +- 每个 Channel Turn 的远端状态只存在于一个 TurnSink。 +- 客户端只渲染服务端快照,不重建领域状态。 +- 中间快照可以丢,最终状态一定可收敛。 +- reasoning 展示、reasoning 回放和“系统正在工作”是三个不同概念。 +- Completed 永远意味着数据库已经提交成功。