808 lines
28 KiB
Markdown
808 lines
28 KiB
Markdown
# 流式 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、正文和工具调用,不能依赖最终字符串中的 `<think>` 后处理。
|
||
|
||
可吸收:
|
||
|
||
- Provider 流的类型化增量;
|
||
- reasoning 与正文分离;
|
||
- 工具参数增量组装;
|
||
- thinking-only 响应的明确处理;
|
||
- reasoning effort/thinking budget 的统一配置概念。
|
||
|
||
### 2.2 ZeroClaw
|
||
|
||
ZeroClaw 把 reasoning 作为不透明 Provider 数据保留,用于要求历史回放的模型,同时处理 `reasoning_content`、`reasoning`、内联 `<think>` 和不同 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<TurnSnapshot>`。生产者覆盖旧值,消费者读取最新值。终态显式编码在快照中,不能仅依赖 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<String>,
|
||
name: Option<String>,
|
||
},
|
||
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<String>,
|
||
pub provider_state: Option<ProviderReasoningState>,
|
||
pub tool_calls: Vec<ToolCall>,
|
||
}
|
||
```
|
||
|
||
约束:
|
||
|
||
- `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<TurnBlock>,
|
||
pub usage: Option<Usage>,
|
||
pub error: Option<String>,
|
||
}
|
||
|
||
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<String>,
|
||
},
|
||
}
|
||
|
||
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<String>,
|
||
},
|
||
}
|
||
```
|
||
|
||
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<ProviderStream, ProviderError>;
|
||
}
|
||
```
|
||
|
||
标题生成、压缩等需要完整响应的代码通过 `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 切分。
|
||
|
||
`<think>`、`<reasoning>` 等内联标签使用有状态 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<Arc<TurnSnapshot>>);
|
||
pub fn begin_finalizing(&mut self);
|
||
pub fn complete(&mut self, usage: Option<Usage>);
|
||
pub fn cancel(&mut self, reason: Option<String>);
|
||
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<TurnSnapshot>`;
|
||
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<Box<dyn TurnSink>, 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(&mut self, snapshot: &TurnSnapshot) -> Result<(), ChannelError>;
|
||
async fn abort(&mut self, snapshot: &TurnSnapshot) -> Result<(), ChannelError>;
|
||
}
|
||
```
|
||
|
||
`TurnSink` 的每个调用都接收完整、过滤后的快照。Sink 不拼接 token。终态方法保留 `&mut self`,使协调器可以在瞬态错误或超时后重试同一个、仍持有远端消息 ID 的 sink;终态成功或重试耗尽后由协调器销毁 sink。
|
||
|
||
### 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 永远意味着数据库已经提交成功。
|