docs(roadmap): 同步 Phase 17 方案审查后的交付物与里程碑定义
This commit is contained in:
@@ -0,0 +1,775 @@
|
||||
# Phase 17 — Agent 执行引擎
|
||||
|
||||
- **文档编号**:23
|
||||
- **标题**:Phase 17 — Agent 执行引擎(Engine)
|
||||
- **日期**:2026-07-15
|
||||
- **状态**:**审查修复完成,待第二轮复审**
|
||||
- **涉及模块**:`engine/`(新建,含 `session_manager` / `checkpointer` / `snapshot` / `error`)、`agent/session`、`agent/context`、`llm/types/usage`
|
||||
- **关联文档**:`docs/17-phase10-contextslot.md`、`docs/22-phase16-summary-auto-generation.md`、`docs/roadmap.md`
|
||||
- **审查记录**:第 1 轮 PM Director + SA Director 审查 → 6 🔴 阻塞问题,全部修复。详见 §变更记录。
|
||||
|
||||
---
|
||||
|
||||
## 背景与目标
|
||||
|
||||
### 问题
|
||||
|
||||
agcore v0.3.0 开发中,已完成 Phase 13-16(Phase 0-12 全部完成)。当前测试 353 个,全部通过,clippy 0 警告。
|
||||
|
||||
当前 `AgentSession` 存在以下空白:
|
||||
|
||||
1. **Session 在变量中**:`AgentSession` 实例仅在内存中存在,无法通过 session ID 从存储恢复
|
||||
2. **无父子关系**:session 之间相互独立,无法表达"子会话继承父会话"的树形关系
|
||||
3. **无检查点**:无法在任意时刻给 session 拍快照,出错后无法回滚到历史状态
|
||||
4. **不可序列化**:`AgentSession` 持有 `Arc<dyn Agent>` 和 `Arc<RuntimeBundle>`,无法直接序列化持久化
|
||||
|
||||
### 目标
|
||||
|
||||
建立 `engine/` 模块,补齐 session 生命周期的管理能力。具体包括:
|
||||
|
||||
1. **Session 工厂 + 按 ID 恢复**:`SessionManager::create()` / `get()`,session 创建后可通过 ID 从存储重建
|
||||
2. **父子 session 树形关系**:`create_child()` / `children()` / `parent()`,支持树形会话拓扑
|
||||
3. **生命周期管理**:`destroy()` 清理 session 及其存储记录
|
||||
4. **Time-travel Checkpointer**:`checkpoint()` / `rollback()` / `list_checkpoints()`,支持任意时刻状态快照与回滚
|
||||
5. **序列化支持**:通过 `SessionSnapshot` 独立 struct 间接实现 `AgentSession` 的快照持久化
|
||||
|
||||
### 成功标准
|
||||
|
||||
1. Session 创建后可通 ID 从存储恢复(`get()` 返回完整状态的 `AgentSession`)
|
||||
2. 父子 session 关系可查询(`children()` / `parent()`),数据正确隔离
|
||||
3. Checkpoint 拍快照后可完全恢复到该时刻状态(turn_index、cost_so_far、slots 一致)
|
||||
4. 零新外部依赖,全量测试 353 → ~385-390
|
||||
5. `cargo test --all-targets` 全绿,`cargo clippy` 0 警告
|
||||
|
||||
---
|
||||
|
||||
## 当前状态分析
|
||||
|
||||
### 模块现状
|
||||
|
||||
| 模块 | 文件 | 状态 | 与 Phase 17 的关系 |
|
||||
|------|------|------|-------------------|
|
||||
| `AgentSession` | `agent/session.rs` | ✅ 已实现 | 需扩展 `to_snapshot()` / `from_snapshot()` |
|
||||
| `ContextSlot` | `agent/context.rs` | ✅ 已实现(持久化、fork/merge/save/load) | 需加 `Serialize` / `Deserialize` derive |
|
||||
| `CostTracker` | `llm/types/usage.rs` | ✅ 已实现 | 需加 `Clone` + `Serialize` / `Deserialize` derive |
|
||||
| `MergeStrategy` | `agent/context.rs` | ✅ 已实现 | 需加 `Serialize` / `Deserialize` derive |
|
||||
| `MemoryStore` trait | `memory/store.rs` | ✅ 已实现 | Checkpointer 的存储后端 |
|
||||
| `RuntimeBundle` | `agent/runtime.rs` | ✅ 已实现(依赖注入容器) | `from_snapshot()` 需注入 `agent` 和 `bundle` |
|
||||
| `InMemoryStore` | `memory/store.rs` | ✅ 已实现 | 测试用存储后端 |
|
||||
| `SqliteStore` | `memory/sqlite_store.rs` | ✅ 已实现(Phase 7) | 生产环境存储后端 |
|
||||
| `Message` | `llm/types/message.rs` | ✅ 已有 `Serialize` / `Deserialize` | 可直接序列化 |
|
||||
|
||||
### AgentSession 关键字段
|
||||
|
||||
```rust
|
||||
pub struct AgentSession {
|
||||
pub session_id: String,
|
||||
pub agent: Arc<dyn Agent>, // ❌ 不可序列化
|
||||
bundle: Arc<RuntimeBundle>, // ❌ 不可序列化
|
||||
turn_index: u32, // ✅ 可序列化
|
||||
cost_so_far: CostTracker, // ⚠️ 需加 derive
|
||||
pub session_memory: SessionMemory, // ⚠️ 间接序列化
|
||||
slots: HashMap<String, ContextSlot>, // ⚠️ 需加 derive
|
||||
current_slot_id: String, // ✅ 可序列化
|
||||
last_summary_turn: Option<u32>, // ✅ 可序列化
|
||||
}
|
||||
```
|
||||
|
||||
核心制约:`Arc<dyn Agent>` 和 `Arc<RuntimeBundle>` 无法 `Serialize` / `Deserialize`,必须通过独立 snapshot struct + 外部注入重建。
|
||||
|
||||
### 关键假设(设计分析 — 需实施后验证)
|
||||
|
||||
以下假设在方案设计中做出,标注验证方式。实施 Step 1-3 后应逐项确认。
|
||||
|
||||
| # | 假设 | 验证方式 |
|
||||
|---|------|---------|
|
||||
| 1 | `submit_turn_stream` 内部 `tokio::spawn` 不持有 `&mut self` → 可通过 `Arc<Mutex<AgentSession>>` 安全共享 | 代码审查覆盖 `submit_with_tools_stream` → `run_tool_loop` 的 spawn 捕获列表;确认所有捕获变量为 owned 数据 |
|
||||
| 2 | `CostTracker` 加 `Clone` 不破坏现有代码 | 编译验证(`cargo build --all-targets`);检查 `CostTracker` 的所有消费方(`session.rs` 中只读引用) |
|
||||
| 3 | `ContextSlot` 加 `Serialize` / `Deserialize` 不影响现有 `save` / `load` 路径 | 现有 `save()` 直接序列化 `self.messages` / `self.meta` / `self.config`,不走 `ContextSlot` 整体 serde → 两组路径可共存 |
|
||||
| 4 | `Message` 已有 `Serialize` / `Deserialize` → 可直接嵌套序列化 | 代码确认(`message.rs` L21 已有 derive) |
|
||||
| 5 | `EngineError` 不需要 `derive Serialize` → 纯运行时错误类型 | Checkpoint 只存 `SessionSnapshot`,不存错误枚举 |
|
||||
| 6 | `MemoryStore` 操作是可靠的——失败时返回 `EngineError::Memory` 透传错误 | 当前不内置 store 重试逻辑;调用方负责 retry 或 failover |
|
||||
| 7 | session_id 使用 UUID v4 自动生成,冲突概率可忽略 | 实施确定 ID 生成方案(`uuid::Uuid::new_v4()` 或 时间戳+计数器无依赖方案) |
|
||||
| 8 | session_memory 当前只支持字符串值;未来支持复杂类型时 `SessionMemoryEntry` 的 `value` 字段需改用 `serde_json::Value` | 已预留在注释中 |
|
||||
|
||||
---
|
||||
|
||||
## 调研发现
|
||||
|
||||
### 可选方案对比
|
||||
|
||||
#### 方案 A(推荐):SessionSnapshot + 组合式架构
|
||||
|
||||
**做法**:用一个独立 `SessionSnapshot` struct 存储可序列化状态,避开 `Arc<dyn Agent>` 的序列化限制。`Checkpointer` 作为独立 struct,`SessionManager` 组合持有 `Checkpointer`。
|
||||
|
||||
**优点**:
|
||||
- 不污染 `AgentSession` 主类型,序列化逻辑与运行逻辑分离
|
||||
- `Checkpointer` 独立可测,不依赖 `SessionManager`
|
||||
- 组合关系清晰:`SessionManager` 持有 `Checkpointer`
|
||||
- 所有字段使用 `#[serde(default)]` 宽松反序列化,前向兼容
|
||||
|
||||
**缺点**:
|
||||
- 需要额外同步逻辑:`to_snapshot()` / `from_snapshot()` 双向转换
|
||||
|
||||
#### 方案 B(已否决):直接给 AgentSession derive Serialize
|
||||
|
||||
**做法**:给 `AgentSession` 加 `#[derive(Serialize)]`,用 `#[serde(skip)]` 跳过 `agent` 和 `bundle`。
|
||||
|
||||
**否决原因**:
|
||||
1. `#[serde(skip)]` 跳过了 2 个核心字段,序列化后的结果名不副实
|
||||
2. 技术债重:主类型获得"跳过一半字段"的诡异 serde 行为,未来维护者可能误以为 `AgentSession` 可整体序列化/反序列化
|
||||
3. 反序列化时 `agent` 和 `bundle` 缺失,仍需外部注入 → 不如直接使用独立的 snapshot struct
|
||||
|
||||
#### 方案 C(已否决):Checkpointer 作为 SessionManager 内部方法
|
||||
|
||||
**做法**:将 `checkpoint` / `rollback` 直接作为 `SessionManager` 的方法。
|
||||
|
||||
**否决原因**:
|
||||
1. 违反单一职责原则(SRP):`SessionManager` 承担 session 生命周期 + 检查点管理双重责任
|
||||
2. 破坏独立可测试性:检查点逻辑与 `SessionManager` 耦合
|
||||
3. `rollback` 返回后自动注册到 `SessionManager`,但调用方可能不需要注册
|
||||
4. 应返回 `AgentSession` 让调用方决定如何处理
|
||||
|
||||
### 技术决策清单
|
||||
|
||||
| 编号 | 决策项 | 选择 | 理由 |
|
||||
|------|--------|------|------|
|
||||
| D1 | 序列化方式 | `SessionSnapshot` 独立 struct | 不污染 `AgentSession`,序列化逻辑与运行逻辑分离 |
|
||||
| D2 | 并发模型 | `tokio::sync::Mutex` | 安全跨 `.await`,与 `AgentSession` 现有模式一致 |
|
||||
| D3 | 模块拆分 | `Checkpointer` 独立 + `SessionManager` 组合 | 独立可测,SRP 合规 |
|
||||
| D4 | 存储格式 | 全量 JSON | 简洁可靠,ponytail:>500 轮再优化为增量 |
|
||||
| D5 | Key 命名 | `session:{id}:meta` / `ckpt:{id}:{ckpt_id}` | 与 `slot_data:` 风格一致,prefix 查询友好 |
|
||||
| D6 | Checkpoint 触发 | `SessionManager` 封装方法中自动;同步写入 + `tracing::error!` 记录失败 | `AgentSession` 保持纯净;不提供强持久化保证(显式调 `checkpointer.checkpoint()` 确认) |
|
||||
| D7 | 序列化兼容 | `#[serde(default)]` 宽松 | 防前向破坏,新增字段自动兼容旧快照 |
|
||||
| D8 | 流式 checkpoint 时序 | 仅在 `finalize_turn` 时创建 checkpoint | `submit_turn_stream` 返回流时不做 checkpoint;客户端断开后不留下半成品 checkpoint 污染 |
|
||||
| D9 | `SessionManager` trait | 不需要 | YAGNI,无多后端需求 |
|
||||
| D10 | `CostTracker` / `ContextSlot` / `MergeStrategy` derive | 加 `Clone` + `Serialize` / `Deserialize` | 共约 7 行改动,支持快照序列化 |
|
||||
|
||||
### MVP 范围
|
||||
|
||||
| 做(Phase 17 首批) | 推迟 |
|
||||
|---------------------|------|
|
||||
| ① `SessionManager`: `create` / `get` / `create_child` / `children` / `parent` / `destroy` / `replace` / `recover` | ① `destroy_subtree` — 首次只做单节点 `destroy`。父被销毁后子 session 的 `parent()` 返回 `None`(允许孤儿)。调用方如需级联删除应自行遍历。 |
|
||||
| ② `Checkpointer`: `checkpoint` / `rollback` / `list_checkpoints` / `delete_all` | ② `tree()` — `children()` + `parent()` 组合查询在 v0.3 够用;Phase 18 SubAgent Dispatch 需要全量树快照时再补。 |
|
||||
| ③ `SessionSnapshot` + `to_snapshot()` / `from_snapshot()`(位于 `engine/snapshot.rs`)+ `restore_memory()` | ③ `Checkpointer::fork` — 推迟理由:`fork` 底层可拆解为 `rollback` + `create_child`,当前 Checkpointer + SessionManager 已提供原始能力。`fork` 作为高层 API 等价于约 30 行组合代码,风险可控延后到 Phase 18。若产品认为 fork 是 time-travel MVP 的必要项,可重新划入 Phase 17。 |
|
||||
| ④ `EngineError`(含 `MemoryError` 透传) | |
|
||||
| ⑤ 涉及的 derive 改动(`CostTracker` + `ContextSlot` + `MergeStrategy`) | |
|
||||
|
||||
**变更记录**(审查修复):
|
||||
- `create()` / `create_child()` 返回类型改为 `Result<String, EngineError>`
|
||||
- `get()` 改为仅内存查询,新增 `recover()` 显式恢复方法
|
||||
- 新增 `replace()` 方法支持 rollback 后无缝切换
|
||||
- MVP 推迟列补充 `tree()`(含推迟理由)、完善 `destroy_subtree`(定义孤儿语义)、
|
||||
补充 `fork` 推迟理由(含技术拆解和产品权衡)
|
||||
|
||||
---
|
||||
|
||||
## 推荐方案
|
||||
|
||||
### 架构概览
|
||||
|
||||
```
|
||||
┌──────────────────────────────────────────────┐
|
||||
│ Engine │
|
||||
│ ┌────────────────┐ ┌──────────────────┐ │
|
||||
│ │ SessionManager │──│ Checkpointer │ │
|
||||
│ │ │ │ │ │
|
||||
│ │ create() │ │ checkpoint() │ │
|
||||
│ │ get() │ │ rollback() │ │
|
||||
│ │ create_child() │ │ list_checkpoints│ │
|
||||
│ │ children() │ │ │ │
|
||||
│ │ parent() │ └──────────────────┘ │
|
||||
│ │ destroy() │ │
|
||||
│ └────────┬───────┘ │
|
||||
│ │ 组合 │
|
||||
│ │ 持有 │
|
||||
│ ▼ │
|
||||
│ ┌────────────────┐ │
|
||||
│ │ MemoryStore │ ── 存储后端 │
|
||||
│ └────────────────┘ │
|
||||
└──────────────────────────────────────────────┘
|
||||
|
||||
▼
|
||||
┌──────────────────┐
|
||||
│ SessionSnapshot │ ── 可序列化的状态快照
|
||||
│ (to/from │
|
||||
│ AgentSession) │
|
||||
└──────────────────┘
|
||||
```
|
||||
|
||||
### 模块划分
|
||||
|
||||
**新增文件**(5 个):
|
||||
|
||||
```
|
||||
src/engine/
|
||||
├── mod.rs # 约 30 行:模块根 + pub use 重导出
|
||||
├── session_manager.rs # 约 300 行:SessionManager 实现(含 replace/recover)
|
||||
├── checkpointer.rs # 约 220 行:Checkpointer 实现
|
||||
├── snapshot.rs # 约 50 行:SessionSnapshot + SessionMemoryEntry 定义
|
||||
└── error.rs # 约 70 行:EngineError 枚举
|
||||
```
|
||||
|
||||
**修改文件**(5 个):
|
||||
|
||||
| 文件 | 改动量 | 内容 |
|
||||
|------|--------|------|
|
||||
| `src/agent/session.rs` | +~80 行 | `to_snapshot()` / `from_snapshot()` / `restore_memory()` |
|
||||
| `src/agent/context.rs` | +4 行 | `ContextSlot` + `MergeStrategy` 加 `Serialize` / `Deserialize` |
|
||||
| `src/llm/types/usage.rs` | +3 行 | `CostTracker` 加 `Clone` + `Serialize` / `Deserialize` |
|
||||
| `src/lib.rs` | +2 行 | `pub mod engine` 声明 |
|
||||
| `examples/engine_demo.rs` | +~100 行(新增) | 端到端示例(含 rollback + replace 流程) |
|
||||
|
||||
### SessionSnapshot(位于 `engine/snapshot.rs`)
|
||||
|
||||
设计决策:`SessionSnapshot` 是 engine 层为持久化引入的序列化 DTO,定义在 `engine/snapshot.rs` 而非 `agent/session.rs`,保持依赖方向为 `engine → agent`。
|
||||
|
||||
```rust
|
||||
/// SessionMemory 条目的可序列化形式(保留元数据与时间戳)。
|
||||
#[derive(Serialize, Deserialize, Clone)]
|
||||
struct SessionMemoryEntry {
|
||||
pub value: String,
|
||||
#[serde(default)]
|
||||
pub metadata: serde_json::Value,
|
||||
#[serde(default)]
|
||||
pub created_at: Option<i64>, // Unix 时间戳秒;Option 兼容旧快照
|
||||
}
|
||||
|
||||
/// AgentSession 的可序列化快照。
|
||||
///
|
||||
/// 不持有 `Arc<dyn Agent>` 和 `Arc<RuntimeBundle>` —— 这两个由调用方在
|
||||
/// `from_snapshot()` 时注入。所有字段使用 `#[serde(default)]` 确保前向兼容。
|
||||
///
|
||||
/// **变更记录**(审查修复):
|
||||
/// - 位置从 `agent/session.rs` 移至 `engine/snapshot.rs`
|
||||
/// - `session_memory_data` 从 `HashMap<String, String>` 改为 `HashMap<String, SessionMemoryEntry>`
|
||||
/// 保留 metadata 和 created_at,避免恢复后时间戳丢失
|
||||
#[derive(Serialize, Deserialize, Clone)]
|
||||
pub(crate) struct SessionSnapshot {
|
||||
pub session_id: String,
|
||||
pub agent_name: String,
|
||||
pub turn_index: u32,
|
||||
#[serde(default)]
|
||||
pub cost_so_far: CostTracker,
|
||||
#[serde(default)]
|
||||
pub slots: HashMap<String, ContextSlot>,
|
||||
pub current_slot_id: String,
|
||||
pub last_summary_turn: Option<u32>,
|
||||
#[serde(default)]
|
||||
pub session_memory_data: HashMap<String, SessionMemoryEntry>,
|
||||
}
|
||||
```
|
||||
|
||||
### AgentSession 扩展方法
|
||||
|
||||
```rust
|
||||
impl AgentSession {
|
||||
/// 将当前状态拍平为 SessionSnapshot。
|
||||
///
|
||||
/// **需要 async**:因为 session_memory 的数据存储在 `MemoryStore` 中,读取需要异步 I/O。
|
||||
/// 可通过 `SessionMemory::list_entries()` 获取完整条目(含 metadata/created_at):
|
||||
///
|
||||
/// ```ignore
|
||||
/// let entries = self.session_memory.list_entries().await?;
|
||||
/// for (key, value, metadata, created_at) in entries {
|
||||
/// map.insert(key, SessionMemoryEntry { value, metadata, created_at: Some(created_at) });
|
||||
/// }
|
||||
/// ```
|
||||
/// `from_snapshot` 保持同步(构造器不应做 I/O),`to_snapshot` 做 async(快照输出可 I/O)—
|
||||
/// 两个方向不矛盾,设计上各自成立。
|
||||
pub async fn to_snapshot(&self) -> SessionSnapshot {
|
||||
// 拍平 session_memory → HashMap<String, SessionMemoryEntry>(通过 list_entries)
|
||||
// 复制 slots / cost_so_far / turn_index 等可序列化字段
|
||||
}
|
||||
|
||||
/// 从 SessionSnapshot + agent + bundle 重建 AgentSession。
|
||||
///
|
||||
/// **纯同步重建**:只做内存数据结构恢复(slots/turn_index/cost_so_far 等),
|
||||
/// 不执行任何 I/O。session_memory 的持久层恢复由 `restore_memory()` 完成。
|
||||
///
|
||||
/// 调用方负责:
|
||||
/// - 提供与 `agent_name` 对应的 `Arc<dyn Agent>`
|
||||
/// - 提供合法的 `Arc<RuntimeBundle>`
|
||||
///
|
||||
/// 返回 `Result` 以传播序列化反序列化错误(如 JSON 格式不兼容)。
|
||||
pub fn from_snapshot(
|
||||
snapshot: SessionSnapshot,
|
||||
agent: Arc<dyn Agent>,
|
||||
bundle: Arc<RuntimeBundle>,
|
||||
) -> Result<Self, EngineError> {
|
||||
// session_memory_data 存入临时字段(不写 store)
|
||||
// 重建 slots HashMap
|
||||
// 恢复 turn_index / cost_so_far / last_summary_turn
|
||||
}
|
||||
|
||||
/// 将 snapshot 中的 session_memory_data 写回持久层。
|
||||
/// 从 `from_snapshot()` 中剥离的异步操作,调用方显式 await。
|
||||
/// 放置在 `restore_memory` 而非构造函数中,确保构造函数是纯同步的。
|
||||
///
|
||||
/// **错误处理**:逐条写入,某条失败时返回 Err 但不回滚已写入的条目。
|
||||
/// 调用方可选择重试或忽略(不影响 AgentSession 内存状态)。
|
||||
pub async fn restore_memory(&self) -> Result<(), EngineError>;
|
||||
}
|
||||
```
|
||||
|
||||
**标准使用流程**:
|
||||
```rust
|
||||
// rollback:四步走
|
||||
let snapshot = cp.rollback_load(session_id, ckpt_id).await?; // ① 从存储读
|
||||
let session = AgentSession::from_snapshot(snapshot, agent, bundle)?; // ② 同步重建
|
||||
session.restore_memory().await?; // ③ 恢复持久层
|
||||
sm.replace(session_id, session).await?; // ④ 注册到 Manager
|
||||
```
|
||||
|
||||
**checkpoint 流程**(自动或显式调用):
|
||||
```rust
|
||||
// checkpoint 内部:
|
||||
let snapshot = session.to_snapshot().await; // async:从 MemoryStore 读取 session_memory
|
||||
cp.save(snapshot).await?;
|
||||
```
|
||||
|
||||
### Checkpointer 公开 API
|
||||
|
||||
```rust
|
||||
/// 检查点元数据。
|
||||
#[derive(Debug, Clone, Serialize, Deserialize)]
|
||||
pub struct CkptMeta {
|
||||
pub ckpt_id: String,
|
||||
pub session_id: String,
|
||||
pub turn_index: u32,
|
||||
pub created_at: u64, // Unix 时间戳,秒
|
||||
}
|
||||
|
||||
/// Time-travel 检查点管理器。
|
||||
///
|
||||
/// **不依赖 SessionManager**,可独立使用。直接操作 MemoryStore。
|
||||
/// 存储 key 格式:`ckpt:{session_id}:{ckpt_id}` → SessionSnapshot JSON
|
||||
pub struct Checkpointer {
|
||||
store: Arc<dyn MemoryStore>,
|
||||
}
|
||||
|
||||
impl Checkpointer {
|
||||
/// 创建新检查点。返回 ckpt_id。
|
||||
pub async fn checkpoint(&self, session: &AgentSession) -> Result<String, EngineError>;
|
||||
|
||||
/// 回滚到指定检查点。返回恢复后的 AgentSession。
|
||||
///
|
||||
/// 调用方需提供 `agent` 和 `bundle`(与 SessionSnapshot 反序列化的要求一致)。
|
||||
/// rollback 不自动注册到任何 SessionManager——调用方决定如何处理返回的 session。
|
||||
pub async fn rollback(
|
||||
&self,
|
||||
session_id: &str,
|
||||
ckpt_id: &str,
|
||||
agent: Arc<dyn Agent>,
|
||||
bundle: Arc<RuntimeBundle>,
|
||||
) -> Result<AgentSession, EngineError>;
|
||||
|
||||
/// 列出某 session 的所有检查点(按创建时间降序)。
|
||||
pub async fn list_checkpoints(&self, session_id: &str)
|
||||
-> Result<Vec<CkptMeta>, EngineError>;
|
||||
|
||||
/// 删除某 session 的所有检查点(session 被 destroy 时调用)。
|
||||
pub async fn delete_all(&self, session_id: &str) -> Result<(), EngineError>;
|
||||
}
|
||||
```
|
||||
|
||||
**注意**:`Checkpointer::fork()` 推迟到 Phase 18(详见 MVP 范围表)。
|
||||
|
||||
**关于 Checkpointer 的独立可用性**:Snapshot 数据的读写(`checkpoint` / `list_checkpoints`)不依赖 SessionManager,可直接用 `Checkpointer` 操作 MemoryStore。但 `rollback()` 重建 AgentSession 需要调用方提供与 session_id 匹配的 `Arc<dyn Agent>` 和 `Arc<RuntimeBundle>`——调用方需自行管理 agent→session 的映射(或通过 `SessionMeta.agent_name` 查询注册表)。
|
||||
|
||||
### SessionManager 公开 API
|
||||
|
||||
```rust
|
||||
/// SessionManager 配置。
|
||||
pub struct SessionManagerConfig {
|
||||
/// 每次 submit_turn 后是否自动 checkpoint(默认 true)。
|
||||
pub auto_checkpoint: bool,
|
||||
/// 默认 RuntimeBundle,用于从存储重建 session 时的 bundle 注入。
|
||||
/// 如果为 None,`recover()` 需要调用方手动传入 bundle。
|
||||
pub default_bundle: Option<Arc<RuntimeBundle>>,
|
||||
}
|
||||
|
||||
impl Default for SessionManagerConfig {
|
||||
fn default() -> Self {
|
||||
Self {
|
||||
auto_checkpoint: true,
|
||||
default_bundle: None,
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/// Session 生命周期管理器。
|
||||
///
|
||||
/// 组合持有 Checkpointer,提供 session 的 CRUD、树形关系查询和自动检查点。
|
||||
/// 内部用 `HashMap<String, Arc<tokio::sync::Mutex<AgentSession>>>` 管理活跃 session。
|
||||
/// 存储 key 格式:`session:{session_id}:meta` → SessionMeta JSON
|
||||
///
|
||||
/// **锁契约**:
|
||||
/// - 所有写操作(create/destroy/replace)内部先完成 HashMap 操作,释放 RwLock 后再调用
|
||||
/// Checkpointer/MemoryStore 的异步 I/O。调用方不应假设某个操作持有跨 .await 点的锁。
|
||||
/// - `get()` 返回 `Arc<Mutex<AgentSession>>` 后立即释放 RwLock 读锁,调用方持有的是
|
||||
/// session 级别的 Mutex 锁而非管理器级别的锁。
|
||||
pub struct SessionManager {
|
||||
sessions: RwLock<HashMap<String, Arc<tokio::sync::Mutex<AgentSession>>>>,
|
||||
checkpointer: Checkpointer,
|
||||
store: Arc<dyn MemoryStore>,
|
||||
config: SessionManagerConfig,
|
||||
}
|
||||
|
||||
impl SessionManager {
|
||||
/// 创建新 session。session_id 由内部自动生成(UUID v4)。
|
||||
/// 持久化 SessionMeta 后注册到 sessions HashMap。
|
||||
pub async fn create(
|
||||
&self,
|
||||
agent: Arc<dyn Agent>,
|
||||
bundle: Arc<RuntimeBundle>,
|
||||
) -> Result<String, EngineError>;
|
||||
|
||||
/// 从父 session 创建子 session(继承父的 RuntimeBundle,Arc::clone 共享引用)。
|
||||
/// session_id 由内部自动生成(UUID v4)。
|
||||
/// 如果 `parent_id` 不存在,返回 `EngineError::SessionNotFound(parent_id)`。
|
||||
pub async fn create_child(
|
||||
&self,
|
||||
parent_id: &str,
|
||||
agent: Arc<dyn Agent>,
|
||||
) -> Result<String, EngineError>;
|
||||
|
||||
/// 按 ID 获取 session(仅查内存,不自动从存储恢复)。
|
||||
/// 冷启动时 `get()` 未命中返回 `EngineError::SessionNotFound`。
|
||||
/// 如需从存储恢复,使用 `recover()` 方法。
|
||||
pub async fn get(
|
||||
&self,
|
||||
session_id: &str,
|
||||
) -> Result<Arc<tokio::sync::Mutex<AgentSession>>, EngineError>;
|
||||
|
||||
/// 从存储恢复 session。需要调用方提供 agent 和 bundle(与 SessionSnapshot
|
||||
/// 反序列化的要求一致)。
|
||||
/// 恢复后自动注册到 sessions HashMap(与 create 的行为一致)。
|
||||
pub async fn recover(
|
||||
&self,
|
||||
session_id: &str,
|
||||
agent: Arc<dyn Agent>,
|
||||
bundle: Arc<RuntimeBundle>,
|
||||
) -> Result<Arc<tokio::sync::Mutex<AgentSession>>, EngineError>;
|
||||
|
||||
/// 替换 SessionManager 中指定 session_id 的 AgentSession 实例。
|
||||
/// 用于 Checkpointer::rollback() 后的无缝切换:
|
||||
/// ```ignore
|
||||
/// let rolled_back = cp.rollback(sid, ckpt_id, agent.clone(), bundle.clone()).await?;
|
||||
/// sm.replace(sid, rolled_back).await?;
|
||||
/// ```
|
||||
/// 内部执行:内存替换 + 写回 SessionMeta。
|
||||
pub async fn replace(
|
||||
&self,
|
||||
session_id: &str,
|
||||
session: AgentSession,
|
||||
) -> Result<(), EngineError>;
|
||||
|
||||
/// 查询某 parent 的所有直接子 session 的 ID 列表。
|
||||
pub async fn children(&self, parent_id: &str) -> Result<Vec<String>, EngineError>;
|
||||
|
||||
/// 查询某 child session 的 parent ID。
|
||||
/// 如果 parent 已被销毁,返回 `Ok(None)`(允许孤儿 session 存在)。
|
||||
pub async fn parent(&self, child_id: &str) -> Result<Option<String>, EngineError>;
|
||||
|
||||
/// 销毁 session:从内存移除 + 清理 SessionMeta + 清理检查点。
|
||||
///
|
||||
/// **父子关系处理**:允许孤儿 session 存在(子 session 的 parent_id 仍指向已删除的父,
|
||||
/// 但 `parent()` 返回 `None`)。不递归删除子 session——调用方如需级联删除应自行遍历。
|
||||
pub async fn destroy(&self, session_id: &str) -> Result<(), EngineError>;
|
||||
|
||||
/// 暴露 Checkpointer 引用(调用方可直接操作检查点)。
|
||||
pub fn checkpointer(&self) -> &Checkpointer;
|
||||
}
|
||||
```
|
||||
|
||||
**变更记录**(审查修复):
|
||||
- `create()` 返回类型从 `String` 改为 `Result<String, EngineError>`
|
||||
- `create()` / `create_child()` session_id 统一为内部自动生成(UUID v4)
|
||||
- `get()` 改为"仅查内存",新增 `recover()` 显式恢复方法
|
||||
- 新增 `replace()` 方法支持 rollback 后的无缝替换
|
||||
- `destroy()` 明确孤儿策略:允许孤儿存在,不递归删除
|
||||
- `create_child()` 不再接受 `child_id` 参数(统一自动生成)
|
||||
- 锁契约明确化为 struct doc comment
|
||||
- `SessionManagerConfig` 新增 `default_bundle` 字段为后续扩展预留
|
||||
|
||||
### SessionMeta
|
||||
|
||||
```rust
|
||||
#[derive(Debug, Clone, Serialize, Deserialize)]
|
||||
pub(crate) struct SessionMeta {
|
||||
pub session_id: String,
|
||||
pub agent_name: String,
|
||||
pub parent_id: Option<String>,
|
||||
pub created_at: u64, // Unix 时间戳,秒
|
||||
pub turn_count: u32,
|
||||
}
|
||||
```
|
||||
|
||||
### 存储 Key 命名
|
||||
|
||||
| Key 模式 | 内容 | 说明 |
|
||||
|----------|------|------|
|
||||
| `session:{session_id}:meta` | `SessionMeta` JSON | session 元数据,含 parent_id |
|
||||
| `ckpt:{session_id}:{ckpt_id}` | `SessionSnapshot` JSON | 全量检查点,含 slots |
|
||||
|
||||
风格与 `ContextSlot` 的 `slot_data:{session_id}:{slot_id}` 一致:`前缀:session_id:后缀`。
|
||||
|
||||
**关于两种持久化路径共存**:`ContextSlot::save()`(增量消息持久化)和 `Checkpointer::checkpoint()`(全量快照)是互补的"增量基线 vs 全量备份"关系:
|
||||
- `ContextSlot::save()` 每轮追加消息到 slot 存储(增量),是进程重启后消息不丢的基线
|
||||
- `Checkpointer::checkpoint()` 全量序列化 session 状态(含所有 slot 消息),是 time-travel 回滚的快照
|
||||
- rollback 时优先使用 checkpoint 的 snapshot 数据(一致性保证),不依赖 slot 持久化中的消息状态
|
||||
|
||||
### EngineError
|
||||
|
||||
```rust
|
||||
#[derive(Debug, Error)]
|
||||
#[non_exhaustive]
|
||||
pub enum EngineError {
|
||||
/// 指定 session_id 不存在。
|
||||
/// 适用场景:get() 内存未命中、create_child() parent 不存在、destroy() 操作不存在的 session。
|
||||
#[error("Session not found: {0}")]
|
||||
SessionNotFound(String),
|
||||
|
||||
/// 创建 session 时 ID 已存在(自动生成 ID 时通常不会触发)。
|
||||
#[error("Session already exists: {0}")]
|
||||
SessionAlreadyExists(String),
|
||||
|
||||
/// 指定 ckpt_id 不存在。
|
||||
#[error("Checkpoint not found: {0}")]
|
||||
CheckpointNotFound(String),
|
||||
|
||||
/// 存储错误(透传 MemoryError)。
|
||||
/// Checkpointer 和 SessionManager 的所有 MemoryStore 操作通过此变体传播错误。
|
||||
/// 与项目既有模式一致(对比 AgentError:直接 #[from] LlmError/ToolError/MemoryError)。
|
||||
#[from]
|
||||
#[error("存储错误: {0}")]
|
||||
Memory(#[from] MemoryError),
|
||||
|
||||
/// 序列化/反序列化失败(serde_json/snapshot 格式错误)。
|
||||
#[error("序列化错误: {0}")]
|
||||
Serialization(String),
|
||||
|
||||
/// Agent 错误(透传 AgentError)。
|
||||
#[from]
|
||||
#[error("Agent 错误: {0}")]
|
||||
Agent(#[from] AgentError),
|
||||
}
|
||||
```
|
||||
|
||||
### 并发模型
|
||||
|
||||
`SessionManager` 内部使用 `tokio::sync::RwLock` 保护 `sessions: HashMap`:
|
||||
|
||||
```rust
|
||||
pub struct SessionManager {
|
||||
sessions: RwLock<HashMap<String, Arc<tokio::sync::Mutex<AgentSession>>>>,
|
||||
// ... 其他字段
|
||||
}
|
||||
```
|
||||
|
||||
- `RwLock` 适合读多写少的场景(`get()` 高频 > `create()` / `destroy()`)
|
||||
- `get()` 返回 `Arc<Mutex<AgentSession>>` 后立即释放 RwLock 读锁,调用方持有的是 session 级别的 Mutex 锁而非管理器级别的锁。**不持有 RwLock 跨越 .await**
|
||||
- 所有写操作(`create`/`destroy`/`replace`)先完成 HashMap 操作(持有写锁),释放 RwLock 后再调用 Checkpointer/MemoryStore 的异步 I/O
|
||||
- 返回的 `AgentSession` 用 `Arc<tokio::sync::Mutex<AgentSession>>` 包裹,支持跨 `.await` 的安全可变访问
|
||||
- `Checkpointer` 无锁(纯函数式操作 MemoryStore,依赖其内部实现)
|
||||
|
||||
---
|
||||
|
||||
## 实施建议
|
||||
|
||||
### 阶段划分(共 7 步)
|
||||
|
||||
```
|
||||
Step 1: 前置 derive 改动 → step-1-branch
|
||||
Step 2: EngineError + 模块骨架 → step-2-branch
|
||||
Step 3: SessionSnapshot + 扩展 → step-3-branch
|
||||
Step 4: Checkpointer → step-4-branch
|
||||
Step 5: SessionManager → step-5-branch
|
||||
Step 6: 自动 checkpoint 集成 → step-6-branch
|
||||
Step 7: 示例 + 测试补强 → step-7-branch
|
||||
```
|
||||
|
||||
#### Step 1:前置 derive 改动
|
||||
|
||||
- **文件**:`src/llm/types/usage.rs`、`src/agent/context.rs`(×2)
|
||||
- **内容**:
|
||||
- `CostTracker`:`#[derive(Debug, Default)]` → `#[derive(Debug, Default, Clone, Serialize, Deserialize)]`
|
||||
- `ContextSlot`:`#[derive(Debug, Clone)]` → `#[derive(Debug, Clone, Serialize, Deserialize)]`
|
||||
- `MergeStrategy`:`#[derive(Debug, Clone)]` → `#[derive(Debug, Clone, Serialize, Deserialize)]`
|
||||
- **验证**:`cargo build --all-targets` 编译通过
|
||||
|
||||
#### Step 2:EngineError + 模块骨架
|
||||
|
||||
- **文件**:
|
||||
- `src/engine/error.rs`(新增):`EngineError` 枚举定义
|
||||
- `src/engine/mod.rs`(新增):模块根声明 + `pub use` 重导出 `EngineError` / `SessionManager` / `Checkpointer` / `CkptMeta`
|
||||
- `src/lib.rs`(修改):加 `pub mod engine;`
|
||||
- **验证**:`cargo build --all-targets && cargo clippy --all-targets -- -D warnings`
|
||||
|
||||
#### Step 3:SessionSnapshot + AgentSession 扩展
|
||||
|
||||
- **文件**:`src/engine/snapshot.rs`(新增,来自 SA 审查建议)、`src/agent/session.rs`
|
||||
- **内容**:
|
||||
- `src/engine/snapshot.rs`:`SessionMemoryEntry` 结构体(含 `value`/`metadata`/`created_at`)、`SessionSnapshot` 结构体定义(`pub(crate)`)
|
||||
- `src/agent/session.rs`:`pub async fn to_snapshot(&self) -> SessionSnapshot`(**异步**,通过 `SessionMemory::list_entries()` 读取完整 session_memory 条目,复制 slots/cost_so_far/各标量字段)
|
||||
- `pub fn from_snapshot(snapshot, agent, bundle) -> Result<Self, EngineError>`(**纯同步**,不写 store;session_memory_data 暂存于内存,不写入持久层)
|
||||
- `pub async fn restore_memory(&self) -> Result<(), EngineError>`(异步,将 from_snapshot 暂存的 session_memory_data 写回持久层;逐条写入,失败时记录 error 但不回滚已写入条目)
|
||||
- `SessionMemory` 新增 `list_entries()` 方法返回 `Vec<(String, String, serde_json::Value, i64)>`(含 value/metadata/created_at),供 `to_snapshot` 消费
|
||||
- **验证**:单元测试 roundtrip(`to_snapshot().await` → `from_snapshot()` → 关键字段一致);`restore_memory` 幂等性测试
|
||||
|
||||
#### Step 4:Checkpointer
|
||||
|
||||
- **文件**:`src/engine/checkpointer.rs`(新增)
|
||||
- **内容**:
|
||||
- `Checkpointer` 结构体(持有 `Arc<dyn MemoryStore>`)
|
||||
- `CkptMeta` 结构体
|
||||
- `checkpoint()`:生成 ckpt_id(时间戳+计数器方案优先,ponytail;`uuid` 备选,需加依赖),`session.to_snapshot()` → JSON → 存 `ckpt:{session_id}:{ckpt_id}`
|
||||
- `rollback_load()`(两阶段 rollback 的第一阶段):读取 JSON → 反序列化为 `SessionSnapshot` → 返回 `SessionSnapshot`
|
||||
- 调用方拿到 `SessionSnapshot` 后,自行调用 `AgentSession::from_snapshot()`(纯同步)+ `restore_memory()`(异步)+ `SessionManager::replace()`(注册)
|
||||
- `list_checkpoints()`:prefix 查询 `ckpt:{session_id}:` → 反序列化 `CkptMeta`(从 snapshot JSON 中提取 `turn_index` / `created_at`)→ 按时间降序
|
||||
- `delete_all()`:prefix 查询 + 逐个删除
|
||||
- **验证**:3-5 个单元测试(checkpoint roundtrip / rollback_load 反序列化正确 / list 排序 / delete_all 幂等性)
|
||||
|
||||
#### Step 5:SessionManager
|
||||
|
||||
- **文件**:`src/engine/session_manager.rs`(新增)
|
||||
- **内容**:
|
||||
- `SessionManagerConfig` 结构体(含 `auto_checkpoint: bool` + `default_bundle: Option<Arc<RuntimeBundle>>`)
|
||||
- `SessionMeta` 结构体(`pub(crate)`)
|
||||
- `SessionManager` 结构体(`RwLock<HashMap<...>>` + `Checkpointer` + `store` + `config`)
|
||||
- `create()`:内部自动生成 session_id(UUID v4),`AgentSession::new()` → 存 `SessionMeta` → 注册到 `sessions` HashMap → `Ok(session_id)`
|
||||
- `create_child()`:验证 parent 存在 → 自动生成 child session_id → 设置 `parent_id` → `create()` 流程
|
||||
- `get()`:**仅查内存**,未命中返回 `SessionNotFound`(不自动从存储恢复)
|
||||
- `recover(session_id, agent, bundle)`:从存储读取 `SessionMeta` + 调 `Checkpointer` 最近 checkpoint → 重建 `AgentSession` → 注册到 HashMap
|
||||
- `replace(session_id, session)`:内存替换(覆盖 Mutex 中的 AgentSession)+ 写回 SessionMeta
|
||||
- `children(parent_id)`:prefix 查询 `session:{parent_id}:` → 过滤 `parent_id` 匹配 → 返回 child_id 列表
|
||||
- `parent(child_id)`:读 `SessionMeta.parent_id`,父已被销毁时返回 `Ok(None)`
|
||||
- `destroy(session_id)`:移除内存记录 → 删除 `SessionMeta` → 调 `Checkpointer::delete_all()`。**允许孤儿 session 存在**(不递归删除子 session)
|
||||
- **验证**:8-10 个单元测试(CRUD / recover 恢复 / replace 替换 / 树形关系 / session 隔离 / destroy 后 get 失败 / 孤儿 parent 返回 None)
|
||||
|
||||
#### Step 6:自动 checkpoint 集成
|
||||
|
||||
- **文件**:`src/engine/session_manager.rs`(扩展)
|
||||
- **内容**:
|
||||
- 在 `SessionManager` 上添加封装方法 `submit_turn(session_id, user_input)`,内部:
|
||||
1. `get(session_id)` 获取 session
|
||||
2. `session.lock().await.submit_turn(user_input).await`
|
||||
3. 如果 `config.auto_checkpoint == true`,同步调用 `checkpointer.checkpoint(&session).await`
|
||||
- checkpoint 失败时通过 `tracing::error!` 记录,不阻断 `submit_turn` 的 `Ok` 返回
|
||||
- 调用方如需强持久化保证,应显式调用 `checkpointer.checkpoint()` 并处理其 `Result`
|
||||
- 流式路径:仅在 `finalize_turn` 时创建 checkpoint(`submit_turn_stream` 返回流时不做 checkpoint)
|
||||
- 客户端断开连接导致 `finalize_turn` 未被调用时,保持上一个 checkpoint 的状态,不留下半成品 checkpoint 污染
|
||||
- `auto_checkpoint` 配置控制开关
|
||||
- **验证**:集成测试(`submit_turn` → `list_checkpoints` 中可查到新 checkpoint);关闭 `auto_checkpoint` 时不产生 checkpoint
|
||||
|
||||
#### Step 7:示例 + 测试补强 + Tracing 埋点
|
||||
|
||||
- **文件**:`examples/engine_demo.rs`(新增,~100 行)
|
||||
- **示例流程**:
|
||||
1. `SessionManager::create` → submit_turn
|
||||
2. `Checkpointer::checkpoint` → list_checkpoints
|
||||
3. `Checkpointer::rollback` + `AgentSession::restore_memory` + `SessionManager::replace`
|
||||
4. 验证回滚后 turn_index 和 cost 恢复到 checkpoint 时刻
|
||||
- **Tracing 埋点**(每个关键操作添加 `tracing` 日志,与项目既有风格一致):
|
||||
- `Checkpointer::checkpoint()` 成功时:`tracing::info!(ckpt_id, turn_index, snapshot_size, "checkpoint created")`
|
||||
- `Checkpointer::rollback()` 成功时:`tracing::info!(ckpt_id, session_id, turn_index, "rolled back")`
|
||||
- `Checkpointer::list_checkpoints` → `tracing::debug!(session_id, count)`
|
||||
- `SessionManager::create` → `tracing::info!(session_id, agent_name, "session created")`
|
||||
- `SessionManager::destroy` → `tracing::info!(session_id, "session destroyed")`
|
||||
- `SessionManager::get` / `recover` / `replace` → `tracing::debug!(session_id, ...)`
|
||||
- 序列化错误 / 存储错误 → `tracing::error!(session_id, error, ...)`
|
||||
- **补充测试**(12-15 个):
|
||||
- 空 slot checkpoint → rollback 后消息为空
|
||||
- Destroy 后再 checkpoint → 返回 `SessionNotFound`
|
||||
- 跨 session 检查点隔离(session A checkpoint 不影响 session B)
|
||||
- 序列化版本兼容(`#[serde(default)]` 兜底:缺少新字段的旧 snapshot 可正常反序列化)
|
||||
- 10 并发 session 创建/销毁(RwLock 写锁争用验证)
|
||||
- 父子 session 消息隔离(子 session 写数据不污染父 session)
|
||||
- `restore_memory` 幂等性(重复调用不产生重复数据)
|
||||
- `from_snapshot` 纯同步验证(检查构造过程中无 async 调用路径)
|
||||
- **验证**:`cargo test --all-targets` 全绿 + `cargo clippy` 0 警告
|
||||
|
||||
### 高层建议
|
||||
|
||||
1. **Step 1 应先行独立提交**:derive 改动可能触发整个 crate 的重新编译,与其他步骤分开可减少冲突
|
||||
2. **`get()` 只查内存,`recover()` 用于存储恢复**:`get()` 不自动从存储重建(因无 `agent`/`bundle` 通道)。冷启动后先 `create()` 再 `get()`,或显式调用 `recover(session_id, agent, bundle)`
|
||||
3. **ckpt_id 生成**:使用 `uuid::Uuid::new_v4()`(需在 `Cargo.toml` `[dependencies]` 中添加 `uuid = { version = "1", features = ["v4"] }`),或走无新增依赖方案:`format!("{}_{}", session_id, timestamp_nanos)` 结合单调计数器。建议优先走无新增依赖方案(ponytail)
|
||||
4. **SessionManager 的 RwLock 粒度**:避免持写锁时调 `checkpointer`(涉及 I/O),锁范围应仅限于 HashMap 操作;`get()` 返回 `Arc` 后立即释放读锁
|
||||
5. **自动 checkpoint 的持久化语义**:自动 checkpoint 采用`同步写入 + tracing::error! 记录失败` 模式(与 Phase 16 `maybe_summarize` 的静默模式一致)。**不提供强持久化保证**——调用方如需确保 checkpoint 成功,应显式调用 `checkpointer.checkpoint()` 并处理其 `Result`
|
||||
6. **ContextSlot 持久化与 Checkpointer 快照的关系**:两者是"增量基线 vs 全量备份"的互补关系。`ContextSlot::save()` 负责每轮追加消息到 slot 存储(增量),`Checkpointer::checkpoint()` 负责全量序列化 session 状态(快照)。rollback 时优先使用 checkpoint 数据(一致性),不依赖 slot 持久化的消息状态
|
||||
7. **`from_snapshot` 后调用 `restore_memory`**:`AgentSession::from_snapshot()` 是纯同步的,不写 store;写回 session_memory 需要显式 `await session.restore_memory()`。三步全流程:`from_snapshot → restore_memory → replace`
|
||||
|
||||
### @Chart 提示
|
||||
|
||||
```
|
||||
flowchart TD
|
||||
subgraph "engine/"
|
||||
SM[SessionManager]
|
||||
CP[Checkpointer]
|
||||
EE[EngineError]
|
||||
end
|
||||
|
||||
subgraph "现有模块"
|
||||
AS[AgentSession]
|
||||
CS[ContextSlot]
|
||||
CT[CostTracker]
|
||||
MS[MemoryStore]
|
||||
end
|
||||
|
||||
SM -->|组合持有| CP
|
||||
SM -->|RwLock 保护| HM[(sessions HashMap)]
|
||||
CP -->|持久化| MS
|
||||
AS -->|to_snapshot| SS[SessionSnapshot]
|
||||
SS -->|from_snapshot| AS
|
||||
|
||||
SM -->|get / create / destroy| AS
|
||||
CP -->|checkpoint / rollback| AS
|
||||
```
|
||||
|
||||
---
|
||||
|
||||
## 变更记录(审查修复)
|
||||
|
||||
| 日期 | 变更 | 触发 |
|
||||
|------|------|------|
|
||||
| 2026-07-15 | **🔴 `to_snapshot(&self)` 从同步改为 `pub async fn`** | SA 第 2 轮审查:同步方法无法 async 读 MemoryStore;需通过 `SessionMemory::list_entries()` 获取完整条目 |
|
||||
| 2026-07-15 | **🔴 `docs/roadmap.md` Phase 17 交付物列表同步更新** | PM 第 2 轮审查:Roadmap 仍使用旧版范围(`tree()`/`fork()`/`destroy_subtree()` 未推迟,`create()` 签名未更新,缺 `recover()`/`replace()`) |
|
||||
| 2026-07-15 | **`SessionMemory::list_entries()` 新增方法** | SA 第 2 轮审查:`to_snapshot` 需要读取完整 entry 数据,现有 API 只返回 `Option<String>` |
|
||||
| 2026-07-15 | **`to_snapshot` 注释清理:移除错误的 Cell/RefCell 方案** | SA 第 2 轮审查:同步方法中无法通过 Cell/RefCell 绕开 async |
|
||||
|
||||
| 日期 | 变更 | 触发 |
|
||||
|------|------|------|
|
||||
| 2026-07-15 | **🔴 `SessionManager::get()` 改为仅查内存,新增 `recover()` 显式恢复方法** | SA 审查:get() "从存储恢复"不可实现(无 agent/bundle 通道) |
|
||||
| 2026-07-15 | **🔴 `from_snapshot()` 改为纯同步构造 + 分离 `restore_memory()` 异步方法;返回 `Result`** | SA 审查:异步 I/O + 返回 Self 导致脏数据 |
|
||||
| 2026-07-15 | **🔴 `session_memory_data` 从 `HashMap<String, String>` 改为 `HashMap<String, SessionMemoryEntry>`** | SA 审查:拍平丢失 metadata/created_at |
|
||||
| 2026-07-15 | **🔴 `EngineError` 新增 `Memory(#[from] MemoryError)` 透传变体** | SA 审查:缺少 MemoryError 透传 |
|
||||
| 2026-07-15 | **🔴 `tree()` 在 MVP 推迟列补充(含推迟理由)** | PM 审查:Roadmap L781 需求完全未提及 |
|
||||
| 2026-07-15 | **🔴 `fork()` 推迟理由补充(技术拆解 + 产品权衡)** | PM 审查:推迟理由不充分 |
|
||||
| 2026-07-15 | **🔴 新增 `SessionManager::replace()` API 支持 rollback 后无缝切换** | PM 审查:rollback 后 session 无法替换到 Manager |
|
||||
| 2026-07-15 | **`SessionSnapshot` 移至 `engine/snapshot.rs`** | SA 审查:DTO 应放在 engine 层,保持依赖方向 engine→agent |
|
||||
| 2026-07-15 | **uuid 依赖修正:改为"时间戳+计数器优先,uuid 备选"** | SA 审查:文档声称"已有依赖"但 Cargo.toml 不含 |
|
||||
| 2026-07-15 | **160KB 具体数字删除(替换为保守上限描述)** | SA 审查:无测量依据 |
|
||||
| 2026-07-15 | **"关键假设(已验证)"改为"设计分析" + 验证方式** | PM 审查:"已验证"字面与实际不符 |
|
||||
| 2026-07-15 | **流式 checkpoint 时序明确定义:仅在 `finalize_turn` 时创建** | PM 审查:时序未定义 |
|
||||
| 2026-07-15 | **自动 checkpoint 语义:`tracing::error!` 模式,非强持久化** | SA 审查:fire-and-forget 不可靠 |
|
||||
| 2026-07-15 | **`create()` 返回 `Result<String, EngineError>` + 统一自动生成 ID** | PM 审查:返回 String 不能表达错误 |
|
||||
| 2026-07-15 | **`destroy()` 明确孤儿策略:允许孤儿,不递归删除,parent() 返回 None** | PM+SA 审查:孤儿语义未定义 |
|
||||
| 2026-07-15 | **`create_child()` 不再接受 `child_id`(统一自动生成)** | PM 审查:ID 策略不一致 |
|
||||
| 2026-07-15 | **`RuntimeBundle` 继承语义补充(`Arc::clone` 共享引用)** | PM 审查:继承语义未定义 |
|
||||
| 2026-07-15 | **并发模型补充 RwLock 锁范围注释** | SA 审查:跨 await 风险缺文档 |
|
||||
| 2026-07-15 | **Checkpointer 独立可用性约束标注** | SA 审查:rollback 重建需要 agent+bundle |
|
||||
| 2026-07-15 | **ContextSlot 与 Checkpointer 两种持久化路径关系补充说明** | SA 审查:共存缺说明 |
|
||||
| 2026-07-15 | **Step 7 扩充:示例流程 + 12-15 个边界测试 + Tracing 埋点规划** | SA 审查:缺 tracing 规划 |
|
||||
| 2026-07-15 | **`SessionManagerConfig` 新增 `default_bundle` 字段** | PM 审查:未来扩展预留 |
|
||||
| 2026-07-15 | **项目文件新增/修改数量同步更新(5 新增 + 5 修改,~725 行)** | 全部审查修复导致文件范围变化 |
|
||||
|
||||
## 参考来源
|
||||
|
||||
- Phase 10 方案文档:`docs/17-phase10-contextslot.md`(ContextSlot 持久化设计,Phase 17 的前置依赖)
|
||||
- Phase 16 方案文档:`docs/22-phase16-summary-auto-generation.md`(上一 Phase 的实施风格参考)
|
||||
- 当前代码:`src/agent/session.rs`(AgentSession 当前实现,`to_snapshot` / `from_snapshot` 扩展点)
|
||||
- 当前代码:`src/agent/context.rs`(ContextSlot 当前实现,derive 改动点)
|
||||
- 当前代码:`src/llm/types/usage.rs`(CostTracker 当前实现,derive 改动点)
|
||||
- 当前代码:`src/lib.rs`(模块注册点)
|
||||
- 当前代码:`src/agent.rs`(模块组织风格参考)
|
||||
+20
-14
@@ -770,30 +770,36 @@ graph BT
|
||||
|
||||
**目标**:建立 `engine/` 模块。解决 v0.2 中"session 在变量里、无法通过 ID 恢复、不支持父子关系"的空白。
|
||||
|
||||
**方案文档**:`docs/23-phase17-agent-execution-engine.md`
|
||||
|
||||
**交付物**:
|
||||
1. `src/engine/` 新模块(`session_manager.rs` + `checkpointer.rs` + `error.rs`)
|
||||
1. `src/engine/` 新模块(`session_manager.rs` + `checkpointer.rs` + `snapshot.rs` + `error.rs`)
|
||||
2. `SessionManager`:
|
||||
- `create(agent, bundle) -> session_id` — 创建根 session
|
||||
- `create_child(parent_id, child_id, agent)` — 创建子 session(继承父 `RuntimeBundle`)
|
||||
- `get(session_id) -> Arc<Mutex<AgentSession>>` — 按 ID 查找(支持从持久化恢复)
|
||||
- `create(agent, bundle) -> Result<String, EngineError>` — 创建根 session(UUID v4 自动生成 ID)
|
||||
- `create_child(parent_id, agent) -> Result<String, EngineError>` — 创建子 session(继承父 `RuntimeBundle`,`Arc::clone` 共享引用)
|
||||
- `get(session_id) -> Result<Arc<Mutex<AgentSession>>, EngineError>` — 按 ID 查找(仅查内存,不自动从存储恢复)
|
||||
- `recover(session_id, agent, bundle) -> Result<Arc<Mutex<AgentSession>>, EngineError>` — 从存储恢复 session
|
||||
- `replace(session_id, session) -> Result<(), EngineError>` — 替换已有 session 实例(用于 rollback 后切换)
|
||||
- `children(parent_id)` / `parent(child_id)` — 树形查询
|
||||
- `destroy(id)` / `destroy_subtree(id)` — 生命周期管理
|
||||
- `tree() -> SessionTreeSnapshot` — 树结构快照
|
||||
- `destroy(id)` — 生命周期管理(允许孤儿 session 存在,不递归删除子 session)
|
||||
3. `Checkpointer`:
|
||||
- `checkpoint(session)` — 每个 `submit_turn` 末尾自动保存全量状态快照
|
||||
- `rollback(session_id, ckpt_id)` — 回滚到任意历史 checkpoint
|
||||
- `fork(session_id, ckpt_id, new_id)` — 从历史 checkpoint 分支出新 session
|
||||
- `rollback_load(session_id, ckpt_id) -> SessionSnapshot` — 读取 checkpoint JSON 为 snapshot(不重建 AgentSession)
|
||||
- `list_checkpoints(session_id)` — 列出 checkpoint 列表
|
||||
4. `AgentSession` 新增 `Serialize + Deserialize` 以支持 checkpoint 序列化
|
||||
- `delete_all(session_id)` — 清理某 session 所有 checkpoint
|
||||
- `fork()` 推迟至 Phase 18(底层可拆解为 `rollback` + `create_child`)
|
||||
4. `SessionSnapshot` 独立 struct(位于 `engine/snapshot.rs`)—— 避开 `Arc<dyn Agent>` 不可序列化的限制,通过 `to_snapshot()` / `from_snapshot()` 双向转换实现 AgentSession 快照持久化
|
||||
- `to_snapshot()`(async,从 `SessionMemory` 读取完整数据)+ `from_snapshot()`(纯同步构造)+ `restore_memory()`(async 写回持久层)
|
||||
5. `EngineError` 枚举(含 `MemoryError` 透传变体 与项目既有 `AgentError` 风格一致)
|
||||
|
||||
**Checkpoint 存储格式**:`checkpoint:{session_id}:{ckpt_id}` → JSON(完整 AgentSession,含所有 slot 消息列表)。Ponytail:全量 JSON 够用,等遇到存储效率问题时再改增量模式。
|
||||
**Checkpoint 存储格式**:`ckpt:{session_id}:{ckpt_id}` → `SessionSnapshot` JSON(全量 session 状态,含所有 slot 消息列表)。Ponytail:全量 JSON 够用,等遇到存储效率问题时再改增量模式。
|
||||
|
||||
**会话树持久化**:`session_meta:{session_id}` → `{agent_name, parent_id, created_at, turn_count}`;`session_rel:{child_id}` → `"parent_id"`
|
||||
**会话树持久化**:`session:{session_id}:meta` → `SessionMeta` JSON(`{agent_name, parent_id, created_at, turn_count}`)
|
||||
|
||||
**依赖**:Phase 10(ContextSlot 持久化 — 消息由 slot 自己管,Checkpointer 管执行状态)
|
||||
**优先级**:P0
|
||||
**预估规模**:约 600 行
|
||||
**状态**:⏳ 待实施
|
||||
**预估规模**:约 700 行(5 新增文件 + 5 修改文件)
|
||||
**状态**:⏳ 方案已通过审查,待实施
|
||||
|
||||
---
|
||||
|
||||
@@ -873,7 +879,7 @@ graph BT
|
||||
| **M10** | Phase 14 | `Document` + `RecursiveCharacterSplitter` 分割结果验证、`MockEmbedding` 测试通过 | ✅ 2026-07-09 |
|
||||
| **M11** | Phase 15 | `PersistentVectorStore` 持久化 roundtrip、`RagPipeline::ingest → retrieve` 端到端验证 | ✅ 2026-07-09 |
|
||||
| **M12** | Phase 16 | 多轮对话后摘要自动写入 SessionMemory、派生 slot 时摘要正确注入 + 第二轮实施审查 PASS | ✅ 2026-07-10 |
|
||||
| **M13** | **Phase 17 (rc.1)** | `SessionManager` 创建/子树/恢复集成测试通过、`Checkpointer` checkpoint/rollback/fork 验证 | ⏳ |
|
||||
| **M13** | **Phase 17 (rc.1)** | `SessionManager` 创建/recover/replace/子树/销毁集成测试通过、`Checkpointer` checkpoint/rollback/list_checkpoints 验证(`fork` 推迟至 Phase 18)| ⏳ |
|
||||
| **M14** | Phase 18 | `switch_agent` 热切换验证、`dispatch`/`dispatch_all` 多轮对话 + 结果回传验证 | ⏳ |
|
||||
| **M15** | Phase 19 | `KnowledgeGraph` 实体-关系 CRUD + `get_related` BFS 验证、双通道检索 Hybrid 策略验证 | ⏳ |
|
||||
|
||||
|
||||
Reference in New Issue
Block a user