docs: design message flow refactor
This commit is contained in:
parent
07b09ca486
commit
65ef919714
@ -2,7 +2,7 @@
|
|||||||
|
|
||||||
本文档描述 PicoBot 当前实现的运行时边界、数据流、并发模型和演进约束。它面向维护者和后续参与改进的 Agent,是代码架构的主入口;行为细节仍以代码和测试为最终依据。
|
本文档描述 PicoBot 当前实现的运行时边界、数据流、并发模型和演进约束。它面向维护者和后续参与改进的 Agent,是代码架构的主入口;行为细节仍以代码和测试为最终依据。
|
||||||
|
|
||||||
流式模型输出、reasoning 展示、活动 Turn 快照和 Channel 实时投递的详细设计与取舍见 [STREAMING_TURN_DESIGN.md](STREAMING_TURN_DESIGN.md)。
|
流式模型输出、reasoning 展示、活动 Turn 快照和 Channel 实时投递的详细设计与取舍见 [STREAMING_TURN_DESIGN.md](STREAMING_TURN_DESIGN.md)。用户输入路由、Session 执行拆分、终态投递确认和历史增量校准的重构方案见 [MESSAGE_FLOW_REFACTOR_DESIGN.md](MESSAGE_FLOW_REFACTOR_DESIGN.md)。
|
||||||
|
|
||||||
## 1. 设计目标
|
## 1. 设计目标
|
||||||
|
|
||||||
|
|||||||
363
docs/MESSAGE_FLOW_REFACTOR_DESIGN.md
Normal file
363
docs/MESSAGE_FLOW_REFACTOR_DESIGN.md
Normal file
@ -0,0 +1,363 @@
|
|||||||
|
# 用户消息到 LLM 回复链路重构设计
|
||||||
|
|
||||||
|
> 状态:实施中(2026-07)。
|
||||||
|
>
|
||||||
|
> 本文定义用户消息入口、Session 执行、Turn 提交和客户端校准链路的重构方案。运行时总览见 `docs/ARCHITECTURE.md`,流式状态模型见 `docs/STREAMING_TURN_DESIGN.md`;代码和测试始终是最终事实来源。
|
||||||
|
|
||||||
|
## 1. 背景
|
||||||
|
|
||||||
|
现有链路已经具备 Channel/Provider 隔离、每 Session 串行、Turn latest-wins 快照和持久化后才发布 `Completed` 等正确基础,但演进过程中留下了以下问题:
|
||||||
|
|
||||||
|
1. `TurnDeliveryService::start` 只报告 sink task 已启动,Session 丢弃 task 的最终结果;sink 终态失败后可能没有普通消息兜底。
|
||||||
|
2. Gateway 用一个循环同步等待所有 inbound 和 control 操作,慢命令或数据库查询会阻塞无关会话。
|
||||||
|
3. Session worker 同时负责上下文准备、Agent 执行、overflow 恢复、提交、投递降级、标题生成和清理。
|
||||||
|
4. 普通工具通知与结构化 `TurnEvent::ToolStarted/ToolFinished` 重复。
|
||||||
|
5. `InboundMessage` 声明了 sender、接收时间和 metadata,但进入 Session 后部分字段被丢弃;平台字段依赖字符串约定透传。
|
||||||
|
6. TUI/WebUI 每次收到终态都重新请求最多 1000 条历史,重复传输刚刚已经通过 TurnSnapshot 下发的结果。
|
||||||
|
7. 自动标题生成在 Turn 完成后仍占用 Session worker,阻塞下一条排队消息。
|
||||||
|
|
||||||
|
## 2. 设计目标
|
||||||
|
|
||||||
|
1. sink 终态失败必须可观测,并至多触发一次普通消息兜底。
|
||||||
|
2. 一个会话的慢 control/command 不得阻塞其他会话的输入和 `/stop`。
|
||||||
|
3. Session worker 只负责队列和生命周期编排,慢步骤由职责明确的 helper/service 承担。
|
||||||
|
4. 上下文首次准备与 overflow 恢复复用同一构建路径。
|
||||||
|
5. 工具进度只有一个权威来源:TurnEvent。
|
||||||
|
6. 用户消息的发送者和接收时间要么被持久化,要么从公共数据契约中删除,不能静默丢失。
|
||||||
|
7. 平台私有上下文以不透明值传递,核心层不解释平台 key。
|
||||||
|
8. 终态提交向客户端提供历史增量;全量历史只用于初次加载、重连和 revision 缺口恢复。
|
||||||
|
9. 标题生成不属于 Turn 完成关键路径,并且迟到结果必须条件提交。
|
||||||
|
10. 所有新增等待、队列、重试和后台任务都必须受 `TaskSupervisor` 管理并有硬边界。
|
||||||
|
|
||||||
|
## 3. 非目标
|
||||||
|
|
||||||
|
- 不合并 `TurnSnapshot` 与持久化 `ChatMessage`;两者分别是暂态展示和耐久事实。
|
||||||
|
- 不把 token delta 放入 MessageBus。
|
||||||
|
- 不取消每 Session 串行语义。
|
||||||
|
- 不让 Channel、Provider 或客户端直接访问 Session 内部状态。
|
||||||
|
- 不在本次重构中改变 SQLite schema 版本;新增消息来源信息复用现有 `source` JSON。
|
||||||
|
- 不保证运行中 Turn 在 Gateway 重启后恢复逐帧状态。
|
||||||
|
|
||||||
|
## 4. 目标数据流
|
||||||
|
|
||||||
|
```text
|
||||||
|
Channel
|
||||||
|
│ normalize + authorize
|
||||||
|
▼
|
||||||
|
InboundEnvelope
|
||||||
|
│
|
||||||
|
▼
|
||||||
|
IngressRouter ───────────────► scoped command task
|
||||||
|
│ │
|
||||||
|
▼ ▼
|
||||||
|
per-session AgentTask queue Session command API
|
||||||
|
│
|
||||||
|
▼
|
||||||
|
ConversationExecutor
|
||||||
|
├── persist user message
|
||||||
|
├── TurnInputBuilder.prepare/recover
|
||||||
|
├── TurnRunner (AgentLoop + cancel)
|
||||||
|
├── TurnCommitter (generation check + atomic persistence)
|
||||||
|
├── DeliveryHandle.await_terminal/fallback
|
||||||
|
└── schedule TitleService
|
||||||
|
│
|
||||||
|
├── TurnSnapshot ─► TurnSink
|
||||||
|
└── TurnCommitted ─► client history delta
|
||||||
|
```
|
||||||
|
|
||||||
|
## 5. 可观测的 Turn 投递
|
||||||
|
|
||||||
|
### 5.1 接口
|
||||||
|
|
||||||
|
`TurnDeliveryService::start` 返回一个必须消费的 handle:
|
||||||
|
|
||||||
|
```rust
|
||||||
|
pub struct TurnDeliveryHandle {
|
||||||
|
completion: oneshot::Receiver<Result<(), DeliveryError>>,
|
||||||
|
}
|
||||||
|
|
||||||
|
impl TurnDeliveryHandle {
|
||||||
|
pub async fn wait(self) -> Result<(), DeliveryError>;
|
||||||
|
}
|
||||||
|
```
|
||||||
|
|
||||||
|
启动失败与异步失败语义分开:
|
||||||
|
|
||||||
|
- `start(...) -> Err`:Channel 不存在、`open_turn` 失败或 supervisor 已停止,Session 从一开始使用普通终态投递。
|
||||||
|
- `start(...) -> Ok(handle)`:sink 生命周期已启动,但不代表终态已经到达外部平台。
|
||||||
|
- `handle.wait() -> Err`:终态重试耗尽或 shutdown abort 失败,Session 执行普通消息兜底。
|
||||||
|
|
||||||
|
### 5.2 兜底规则
|
||||||
|
|
||||||
|
1. Agent 结果必须先持久化,之后 Turn 才能 `Completed`。
|
||||||
|
2. Session 发布终态后等待 delivery handle;等待本身由 coordinator 的 sink timeout/retry 限制。
|
||||||
|
3. delivery 成功:不发送普通消息。
|
||||||
|
4. delivery 失败:通过 `MessageBus::deliver_outbound` 发送最终正文,并记录明确错误。
|
||||||
|
5. `Cancelled`/`Failed` 终态不重复发送正文;只有存在可展示 partial 且 sink 失败时才发送 partial/failure 摘要。
|
||||||
|
6. Channel sink 内部可以做平台特定编辑降级,但不得把“未找到目标/未发送”报告为成功。
|
||||||
|
|
||||||
|
### 5.3 测试
|
||||||
|
|
||||||
|
- open 失败时发送一次普通终态。
|
||||||
|
- open 成功、finish 永久失败时发送一次普通终态。
|
||||||
|
- finish 瞬态失败后成功时不发送普通终态。
|
||||||
|
- 持久化失败时不得把失败前正文作为 completed fallback 发送。
|
||||||
|
- shutdown/cancel 不造成双重终态。
|
||||||
|
|
||||||
|
## 6. Gateway 入口并发
|
||||||
|
|
||||||
|
### 6.1 现状问题
|
||||||
|
|
||||||
|
单个 `message-processor` 在 `tokio::select!` 分支内等待 `handle_message` 和 control I/O。`/compact` 的 LLM 调用、大历史查询或慢 SQLite 操作会造成跨 Session 队头阻塞。
|
||||||
|
|
||||||
|
### 6.2 目标
|
||||||
|
|
||||||
|
拆成两个只负责消费和派发的 supervisor task:
|
||||||
|
|
||||||
|
- `inbound-router`:消费 `InboundMessage`,为每条输入启动受监督的短派发任务;普通消息最终进入 Session 队列。
|
||||||
|
- `control-router`:消费 `ControlMessage`,为每个请求启动受监督任务并通过一次性回复通道返回。
|
||||||
|
|
||||||
|
Session 内部继续用 mutex、`persistence_lock`、`worker_generation` 和 `state_version` 保证同一对话的一致性。Gateway 不再通过全局串行获得隐式正确性。
|
||||||
|
|
||||||
|
### 6.3 有界性
|
||||||
|
|
||||||
|
- MessageBus 仍是全局 admission queue。
|
||||||
|
- 普通 AgentTask 仍受每 Session 容量 32 限制。
|
||||||
|
- router 通过 `TaskSupervisor::spawn` 管理请求任务;spawn 失败必须向调用者/Channel 返回错误。
|
||||||
|
- control reply 使用 `oneshot`,每个请求只有一个结果。
|
||||||
|
- `/stop` 直接修改目标 Session cancellation/generation,不进入 AgentTask 队列。
|
||||||
|
|
||||||
|
### 6.4 测试
|
||||||
|
|
||||||
|
- Session A 的慢 command 不阻塞 Session B 的普通输入。
|
||||||
|
- Session A 的慢 history query 不阻塞 Session B `/stop`。
|
||||||
|
- router shutdown 后新请求收到明确失败。
|
||||||
|
- 同一 Session 的 AgentTask 顺序保持不变。
|
||||||
|
|
||||||
|
## 7. Session 执行拆分
|
||||||
|
|
||||||
|
### 7.1 `TurnInputBuilder`
|
||||||
|
|
||||||
|
输入:稳定的 `SessionTurnSnapshot`、用户输入、skills、MemoryManager、WorkManager。
|
||||||
|
|
||||||
|
输出:
|
||||||
|
|
||||||
|
```rust
|
||||||
|
struct PreparedTurnInput {
|
||||||
|
messages: Vec<ChatMessage>,
|
||||||
|
base_state_version: u64,
|
||||||
|
compression_update: Option<CompressionUpdate>,
|
||||||
|
}
|
||||||
|
```
|
||||||
|
|
||||||
|
职责:
|
||||||
|
|
||||||
|
- 并发读取 Knowledge memory 和 active plan;
|
||||||
|
- 运行 ContextCompressor;
|
||||||
|
- 统一插入 system prompt;
|
||||||
|
- 统一向最后一条用户消息追加 runtime context;
|
||||||
|
- 返回需要条件提交的 compression metadata,不直接持有 Session 锁做慢 I/O。
|
||||||
|
|
||||||
|
`recover_after_overflow` 复用同一 assembly 函数,只替换 context window 和压缩结果,不能复制 prompt/runtime context 拼装逻辑。
|
||||||
|
|
||||||
|
### 7.2 `TurnRunner`
|
||||||
|
|
||||||
|
职责:
|
||||||
|
|
||||||
|
- 创建 `TurnController`、`AgentTurnContext` 和 delivery handle;
|
||||||
|
- 在 `AgentLoop` 与 cancel receiver 之间 select;
|
||||||
|
- 最多执行一次 context-overflow recovery;
|
||||||
|
- 返回类型化 `TurnRunOutcome`,不直接构造 OutboundMessage。
|
||||||
|
|
||||||
|
```rust
|
||||||
|
enum TurnRunOutcome {
|
||||||
|
Completed(AgentProcessResult),
|
||||||
|
Cancelled(Option<ChatMessage>),
|
||||||
|
Failed { error: AgentError, partial: Option<ChatMessage> },
|
||||||
|
Stale,
|
||||||
|
}
|
||||||
|
```
|
||||||
|
|
||||||
|
### 7.3 `TurnCommitter`
|
||||||
|
|
||||||
|
职责:
|
||||||
|
|
||||||
|
- 提交前验证 generation/state version;
|
||||||
|
- 原子持久化 `emitted_messages`;
|
||||||
|
- 成功后发布 `Completed`;
|
||||||
|
- 失败/取消时按 partial 规则持久化并发布对应终态;
|
||||||
|
- 产生 `CommittedTurnDelta`。
|
||||||
|
|
||||||
|
### 7.4 Session worker 保留职责
|
||||||
|
|
||||||
|
- 从队列接收 AgentTask;
|
||||||
|
- 持久化原始用户消息;
|
||||||
|
- 捕获稳定快照;
|
||||||
|
- 顺序调用 builder/runner/committer;
|
||||||
|
- 清理 active turn 和 cancel handle;
|
||||||
|
- 调度非关键后台工作。
|
||||||
|
|
||||||
|
## 8. 上下文策略收口
|
||||||
|
|
||||||
|
上下文策略分两级,但 owner 明确:
|
||||||
|
|
||||||
|
- `TurnInputBuilder`:跨轮历史压缩、Timeline/Memory/Plan、overflow recovery。
|
||||||
|
- `AgentLoop`:单次工具循环中临时裁剪过大的旧 tool result,不修改 Session 历史。
|
||||||
|
|
||||||
|
二者不能重复构建 system/runtime prompt。`AgentLoop` 的“缺 system 时自动注入”仅保留给明确的 stateless API;交互 Session 调用使用要求首条必须为 system 的入口或 debug assertion。
|
||||||
|
|
||||||
|
Memory recall 和 active plan 查询互不依赖,应使用 `tokio::join!` 并发执行。任何结果提交前都验证 `base_state_version`。
|
||||||
|
|
||||||
|
## 9. 工具进度唯一来源
|
||||||
|
|
||||||
|
交互 Turn 删除 `AgentLoop.notify_tx: UnboundedSender<String>` 和每消息 notification publisher。工具进度仅由:
|
||||||
|
|
||||||
|
```text
|
||||||
|
ToolStarted → TurnController → TurnSnapshot
|
||||||
|
ToolFinished → TurnController → TurnSnapshot
|
||||||
|
```
|
||||||
|
|
||||||
|
后台子 Agent 的 `TaskNotification` 是另一种领域事件,继续保留,因为它表达跨 Turn 的任务完成,而不是当前 Turn 的工具进度。
|
||||||
|
|
||||||
|
## 10. Inbound 与 ChannelContext
|
||||||
|
|
||||||
|
### 10.1 规范化输入
|
||||||
|
|
||||||
|
```rust
|
||||||
|
struct InboundMessage {
|
||||||
|
channel: String,
|
||||||
|
chat_id: String,
|
||||||
|
sender_id: String,
|
||||||
|
content: String,
|
||||||
|
received_at: i64,
|
||||||
|
media: Vec<MediaItem>,
|
||||||
|
channel_context: ChannelContext,
|
||||||
|
}
|
||||||
|
```
|
||||||
|
|
||||||
|
`metadata` 与 `forwarded_metadata` 合并为语义明确的不透明 `ChannelContext`。核心只允许:
|
||||||
|
|
||||||
|
- 原样传给本轮 `TurnTarget` 或普通错误回复;
|
||||||
|
- 从通用 typed 字段读取 `reply_to`;
|
||||||
|
- 不解析 `feishu.*` 等平台 key。
|
||||||
|
|
||||||
|
平台 message/reaction ID 最终由具体 sink 持有。`feishu.parent_id` 要么映射为 typed `reply_to`,要么删除,不能继续作为无消费者字段。
|
||||||
|
|
||||||
|
### 10.2 持久化用户来源
|
||||||
|
|
||||||
|
用户 `ChatMessage` 使用原始 `received_at`,并设置:
|
||||||
|
|
||||||
|
```rust
|
||||||
|
MessageSource {
|
||||||
|
kind: UserInput,
|
||||||
|
from_channel: Some(channel),
|
||||||
|
from_user_id: Some(sender_id),
|
||||||
|
...
|
||||||
|
}
|
||||||
|
```
|
||||||
|
|
||||||
|
客户端历史投影可以隐藏内部 sender ID;LLM 上下文是否展示发言者由独立策略决定,不能直接泄露平台标识。
|
||||||
|
|
||||||
|
## 11. 终态历史增量
|
||||||
|
|
||||||
|
### 11.1 协议
|
||||||
|
|
||||||
|
持久化成功后发布:
|
||||||
|
|
||||||
|
```rust
|
||||||
|
WsOutbound::TurnCommitted {
|
||||||
|
session_id: String,
|
||||||
|
history_revision: u64,
|
||||||
|
messages: Vec<HistoryMessage>,
|
||||||
|
}
|
||||||
|
```
|
||||||
|
|
||||||
|
规则:
|
||||||
|
|
||||||
|
- `TurnUpdated(Completed)` 仍负责结束 active turn 展示。
|
||||||
|
- `TurnCommitted` 只包含本轮新持久化消息,负责把 transcript 校准到数据库事实。
|
||||||
|
- Session 维护单调 `history_revision`;客户端只接受连续 revision。
|
||||||
|
- revision 缺口、重连或显式切换 Session 时才请求全量历史。
|
||||||
|
- 全量 `SessionHistory` 返回当前 revision。
|
||||||
|
|
||||||
|
第一阶段可使用最终 assistant/tool message IDs 去重;如果不修改数据库 schema,revision 使用 Session 内 `state_version`/最新 message sequence 投影,重启后从 Storage 最大 sequence 恢复。
|
||||||
|
|
||||||
|
### 11.2 测试
|
||||||
|
|
||||||
|
- 正常 Turn 完成不触发全量 history 请求。
|
||||||
|
- 增量包含 assistant tool call、tool result 和最终 assistant。
|
||||||
|
- 重复增量按 message ID 幂等。
|
||||||
|
- revision 缺口触发一次全量校准。
|
||||||
|
- 其他 Session 的增量只更新对应缓存/未读状态。
|
||||||
|
|
||||||
|
## 12. 标题后台化
|
||||||
|
|
||||||
|
Turn 提交后,Session worker 调用 `TitleService::schedule` 并立即处理下一条任务。
|
||||||
|
|
||||||
|
后台任务:
|
||||||
|
|
||||||
|
1. 在锁内捕获 title prompt、session ID 和 `state_version`。
|
||||||
|
2. 在锁外调用 Provider。
|
||||||
|
3. 获取 `persistence_lock`。
|
||||||
|
4. 只有标题仍为默认值且 generation/state 条件允许时提交。
|
||||||
|
5. 任务由 `TaskSupervisor` 管理;shutdown 时取消并限时回收。
|
||||||
|
|
||||||
|
标题失败只记录 warning,不改变 Turn 状态,不向用户发送错误消息。
|
||||||
|
|
||||||
|
## 13. 错误模型
|
||||||
|
|
||||||
|
Session worker 不再散落构造英文字符串 OutboundMessage,而是返回类型化错误:
|
||||||
|
|
||||||
|
```rust
|
||||||
|
enum TurnFailureKind {
|
||||||
|
InputPersistence,
|
||||||
|
AgentCreation,
|
||||||
|
ContextPreparation,
|
||||||
|
Provider,
|
||||||
|
TurnPersistence,
|
||||||
|
Delivery,
|
||||||
|
}
|
||||||
|
```
|
||||||
|
|
||||||
|
统一 `TurnFailurePresenter` 根据 Channel/PresentationPolicy 生成用户可见内容。原始 provider/storage 错误只进入安全日志和 Turn internal error,不直接暴露 secrets。
|
||||||
|
|
||||||
|
Gateway 的 `handle_message` 错误不能只写日志;必须通过输入携带的 ChannelContext 返回一个关联到原消息的错误结果。
|
||||||
|
|
||||||
|
## 14. 迁移与提交顺序
|
||||||
|
|
||||||
|
1. 文档:落地本设计和架构链接。
|
||||||
|
2. Delivery:返回 completion handle,Session 消费结果并做一次兜底。
|
||||||
|
3. Gateway:拆分 inbound/control router,消除全局慢操作串行。
|
||||||
|
4. Session:提取输入构建和 overflow recovery,再提取 run/commit helper。
|
||||||
|
5. Progress/title:删除旧工具通知,标题移入 supervisor。
|
||||||
|
6. Contract:规范 InboundMessage、用户来源和 ChannelContext。
|
||||||
|
7. Protocol:增加 TurnCommitted/history revision,客户端改为增量校准。
|
||||||
|
8. 最终清理:删除不可达 `HandleResult::AgentResponse` 交互分支和重复 helper。
|
||||||
|
|
||||||
|
每一步都保持可编译、可测试、可单独回滚,不允许一个提交同时更改全部并发和协议语义。
|
||||||
|
|
||||||
|
## 15. 回归测试矩阵
|
||||||
|
|
||||||
|
| 风险 | 必需测试 |
|
||||||
|
|---|---|
|
||||||
|
| sink 异步终态失败 | fallback 恰好一次,成功时零次 |
|
||||||
|
| Gateway 队头阻塞 | 慢 A 不阻塞 B 输入/control |
|
||||||
|
| stale worker | generation/state 改变后不得提交 |
|
||||||
|
| overflow recovery | prompt/runtime context 只附加一次 |
|
||||||
|
| duplicate tool progress | 交互 Turn 不产生普通工具通知 |
|
||||||
|
| inbound fidelity | sender、received_at、reply_to 正确保留 |
|
||||||
|
| title race | 用户重命名后迟到标题不得覆盖 |
|
||||||
|
| history delta | 连续、重复、缺口、跨 Session |
|
||||||
|
| shutdown | router、delivery、title task 均被监督和有界回收 |
|
||||||
|
|
||||||
|
## 16. 完成条件
|
||||||
|
|
||||||
|
- `cargo test --lib`
|
||||||
|
- `cargo test --test test_scheduler --test test_request_format`
|
||||||
|
- `cargo clippy --all-targets --all-features -- -D warnings`
|
||||||
|
- `cargo build`
|
||||||
|
- `webui/npm run check`
|
||||||
|
- `webui/npm run build`
|
||||||
|
- 新增的失败、取消、并发、revision 和 fallback 测试全部通过。
|
||||||
|
- `docs/ARCHITECTURE.md`、AGENTS.md 中的运行时不变量与实现一致。
|
||||||
|
|
||||||
Loading…
x
Reference in New Issue
Block a user