PicoBot/docs/MESSAGE_FLOW_REFACTOR_DESIGN.md

14 KiB
Raw Permalink Blame History

用户消息到 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. 目标数据流

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

pub struct TurnDeliveryHandle {
    completion: oneshot::Receiver<Result<(), DeliveryError>>,
}

impl TurnDeliveryHandle {
    pub async fn wait(self) -> Result<(), DeliveryError>;
}

启动失败与异步失败语义分开:

  • start(...) -> ErrChannel 不存在、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-processortokio::select! 分支内等待 handle_message 和 control I/O。/compact 的 LLM 调用、大历史查询或慢 SQLite 操作会造成跨 Session 队头阻塞。

6.2 目标

拆成两个只负责消费和派发的 supervisor task

  • inbound-router:消费 InboundMessage,为每条输入启动受监督的短派发任务;普通消息最终进入 Session 队列。
  • control-router:消费 ControlMessage,为每个请求启动受监督任务并通过一次性回复通道返回。

Session 内部继续用 mutex、persistence_lockworker_generationstate_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。

输出:

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

职责:

  • 创建 TurnControllerAgentTurnContext 和 delivery handle
  • AgentLoop 与 cancel receiver 之间 select
  • 最多执行一次 context-overflow recovery
  • 返回类型化 TurnRunOutcome,不直接构造 OutboundMessage。
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。工具进度仅由

ToolStarted → TurnController → TurnSnapshot
ToolFinished → TurnController → TurnSnapshot

后台子 Agent 的 TaskNotification 是另一种领域事件,继续保留,因为它表达跨 Turn 的任务完成,而不是当前 Turn 的工具进度。

10. Inbound 与 ChannelContext

10.1 规范化输入

struct InboundMessage {
    channel: String,
    chat_id: String,
    sender_id: String,
    content: String,
    received_at: i64,
    media: Vec<MediaItem>,
    channel_context: ChannelContext,
}

metadataforwarded_metadata 合并为语义明确的不透明 ChannelContext。核心只允许:

  • 原样传给本轮 TurnTarget 或普通错误回复;
  • 从通用 typed 字段读取 reply_to
  • 不解析 feishu.* 等平台 key。

平台 message/reaction ID 最终由具体 sink 持有。feishu.parent_id 要么映射为 typed reply_to,要么删除,不能继续作为无消费者字段。

10.2 持久化用户来源

用户 ChatMessage 使用原始 received_at,并设置:

MessageSource {
    kind: UserInput,
    from_channel: Some(channel),
    from_user_id: Some(sender_id),
    ...
}

客户端历史投影可以隐藏内部 sender IDLLM 上下文是否展示发言者由独立策略决定,不能直接泄露平台标识。

11. 终态历史增量

11.1 协议

持久化成功后发布:

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 去重;如果不修改数据库 schemarevision 使用 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而是返回类型化错误

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 handleSession 消费结果并做一次兜底。
  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 中的运行时不变量与实现一致。