Files
DevFlow/docs/02-架构设计/专项设计/消息拆分存储设计-2026-06-19.md
绝尘 44d1c6a00c 优化: 文档路径同步 + 审查登记 + 消息级溯源设计文档
- PROGRESS/URGENT 文档分类后子目录路径修正
- 待审查.md 登记 CR-260619-08 intent + CR-260619-09 PhaseB+C(均  PASS)
- 新增消息级溯源设计(消息拆分存储 + 知识/审计/灵感全链路消息级定位)
2026-06-19 18:58:49 +08:00

309 lines
14 KiB
Markdown

# 消息拆分存储设计
> 创建:2026-06-19 | 编号:F-260619-03(消息拆分) | 状态:📐 设计待实施
> 关联任务:`消息拆分存储:ai_messages 表 + 全量迁移 + 三阶段渐进切换`(DevFlow 项目)
> 关联灵感:b2e61a21 消息拆分存储 + 消息级溯源 + 知识脉络网络
> 关联设计:[消息级溯源设计-2026-06-19.md](./消息级溯源设计-2026-06-19.md)(本设计的前置依赖 + 消费方)
> 上级索引:[../INDEX.md](../INDEX.md)
---
## 一、目标
`ai_conversations.messages`(整个对话 JSON 数组存于单 TEXT 列)拆分为独立的 `ai_messages` 表(每条消息一行),实现:
- **增量写入**:长对话(50+ 轮)save 不再全量序列化覆盖写
- **O(1) 删除/压缩/单条编辑**:clear/compress/replace_tool_result_content 改为 SQL 单行操作
- **游标分页加载**:超长对话按 seq 范围拉取(未来)
- **持久化消息 ID**:为[消息级溯源](./消息级溯源设计-2026-06-19.md)提供 O(1) 查询基础
---
## 二、现状
### 2.1 存储模型
```
ai_conversations.messages TEXT -- 存整个 Vec<ChatMessage> 序列化 JSON
```
- `save_conversation()`(`conversation.rs:135`)全量序列化覆盖写
-`chat.rs` / `agentic.rs` / 审批 / 压缩 / 编辑重生成等 **10+ 处**调用
- `ContextManager`(`crates/df-ai/src/context.rs`)的 push/pop/replace_tool_result_content/compress_old_messages/insert_at/restore_from_messages 全部基于"内存 Vec"假设
### 2.2 ChatMessage 结构(当前 9 字段)
定义在 `crates/df-ai-core/src/types.rs:72`:
| 字段 | 类型 | 说明 |
|---|---|---|
| `role` | `MessageRole` | system/user/assistant/tool |
| `content` | `String` | 文本内容(parts 的 Text 片与 content 一致) |
| `parts` | `Option<Vec<ContentPart>>` | 多模态片(F-05 Phase 2a) |
| `tool_call_id` | `Option<String>` | role=Tool 时必填 |
| `tool_calls` | `Option<Vec<ToolCall>>` | AI 发起的工具调用 |
| `model` | `Option<String>` | 生成该消息的 model(仅 assistant) |
| `status` | `Option<String>` | None/"active" 正常;"truncated" 软删 |
| `reasoning_content` | `Option<String>` | DeepSeek thinking 推理内容 |
| `timestamp` | `Option<i64>` | Unix 毫秒(前端展示用) |
> **本设计新增**:`id: Option<String>`(ULID)——见[消息级溯源设计](./消息级溯源设计-2026-06-19.md) §三。
### 2.3 迁移基线
- 当前最新迁移版本:**V19**(`migrations.rs:45`)
- V20 预留给 F-260619-01(任务关联灵感)
- **本设计使用 V21**(建 ai_messages 表 + 全量迁移)
- V22 预留给 Phase 3 删旧列
---
## 三、数据模型
### 3.1 新建 `ai_messages` 表
```sql
CREATE TABLE IF NOT EXISTS ai_messages (
id TEXT PRIMARY KEY, -- ULID 全局唯一 + 时间有序
conversation_id TEXT NOT NULL, -- 关联对话
seq INTEGER NOT NULL, -- 对话内序号(排序用,从 0 起)
role TEXT NOT NULL, -- system/user/assistant/tool
content TEXT NOT NULL DEFAULT '',
parts TEXT, -- 多模态 parts JSON(可空)
tool_call_id TEXT,
tool_calls TEXT, -- AI 发起的工具调用 JSON(可空)
model TEXT, -- 生成该消息的 model(仅 assistant)
status TEXT NOT NULL DEFAULT 'active', -- active/compressed/archived_segment/truncated
reasoning_content TEXT,
timestamp INTEGER, -- 消息创建时间(Unix 毫秒)
created_at TEXT NOT NULL,
UNIQUE(conversation_id, seq)
);
CREATE INDEX IF NOT EXISTS idx_ai_messages_conv ON ai_messages(conversation_id, seq);
```
### 3.2 `ai_conversations` 表保留 messages 列
- Phase 2 双写期:旧列做回退保险
- Phase 3 切读稳定后:V22 迁移删列
---
## 四、渐进路径(三阶段)
### Phase 1:脏标记增量写入(不改 schema)
**目标**:减少 save_conversation 的序列化开销。
- `ContextManager``dirty_min_seq: Option<usize>` / `dirty_max_seq: Option<usize>` 标记变化范围
- `save_conversation` 只序列化变化范围内的消息(但仍写整列)
- push/pop/replace/compress/insert_at 各操作设置 dirty 标记
- save 后清空 dirty 标记
**风险**:低,纯内存优化,不涉及 schema 变更。
### Phase 2:建表 + 双写 + 全量迁移
**目标**:数据落入新表,读仍走旧列。
#### 4.2.1 迁移 V21
```rust
fn migrate_v21(conn: &Connection) -> Result<()> {
// 1. 建表(IF NOT EXISTS 幂等)
conn.execute_batch(V21_SQL)?;
// 2. COUNT 探测:ai_messages 已有数据 → 跳过迁移只写版本号
let existing: i64 = conn.query_row(
"SELECT COUNT(*) FROM ai_messages", [], |row| row.get(0)
)?;
if existing > 0 {
tracing::info!("v21: ai_messages 已有 {} 条,跳过迁移", existing);
conn.execute("INSERT INTO schema_version (version) VALUES (?)", [21])?;
return Ok(());
}
// 3. 遍历 ai_conversations
let mut stmt = conn.prepare("SELECT id, messages, created_at FROM ai_conversations")?;
let rows = stmt.query_map([], |row| {
Ok((row.get::<_, String>(0)?, row.get::<_, String>(1)?, row.get::<_, String>(2)?))
})?;
let all_rows: Vec<_> = rows.collect::<Result<Vec<_>, _>>()?;
// 4. 分批 commit(每 50 个对话一批,避免长事务持有写锁)
const BATCH_SIZE: usize = 50;
for (batch_idx, batch) in all_rows.chunks(BATCH_SIZE).enumerate() {
let tx = conn.unchecked_transaction()?;
for (conv_id, messages_json, conv_created_at) in batch {
// 5. 逐对话反序列化 messages JSON → Vec<serde_json::Value>
// (用裸 JSON 而非 ChatMessage,因 df-storage 不依赖 df-ai-core)
let messages: Vec<serde_json::Value> = match serde_json::from_str(messages_json) {
Ok(v) => v,
Err(e) => {
tracing::warn!("v21: 对话 {} messages JSON 解析失败,跳过: {}", conv_id, e);
continue; // 坏数据跳过,不中断迁移
}
};
for (seq, msg) in messages.iter().enumerate() {
// 6. 逐条消息提取字段 → INSERT INTO ai_messages
// 字段名硬编码("role"/"content" 等)——ChatMessage 改名会漏数据!
// 耦合点标注:见下方「迁移耦合点」
let id = format!("msg_migrated_{}_{}", conv_id, seq); // 迁移期 ID,天然唯一
let role = msg.get("role").and_then(|v| v.as_str()).unwrap_or("user");
let content = msg.get("content").and_then(|v| v.as_str()).unwrap_or("");
let parts = msg.get("parts")
.filter(|v| !v.is_null())
.map(|v| v.to_string());
let tool_call_id = msg.get("tool_call_id")
.and_then(|v| v.as_str())
.map(String::from);
let tool_calls = msg.get("tool_calls")
.filter(|v| !v.is_null())
.map(|v| v.to_string());
let model = msg.get("model")
.and_then(|v| v.as_str())
.map(String::from);
// status 归一化:None/空 → "active"(列语义清晰,永不 NULL)
let status = msg.get("status")
.and_then(|v| v.as_str())
.filter(|s| !s.is_empty())
.unwrap_or("active");
let reasoning_content = msg.get("reasoning_content")
.and_then(|v| v.as_str())
.map(String::from);
// created_at:有 timestamp 用消息自己的,没有 fallback 到对话创建时间
let timestamp = msg.get("timestamp").and_then(|v| v.as_i64());
let created_at = timestamp
.map(|ts| ts.to_string())
.unwrap_or_else(|| conv_created_at.clone());
tx.execute(
"INSERT OR IGNORE INTO ai_messages
(id, conversation_id, seq, role, content, parts, tool_call_id,
tool_calls, model, status, reasoning_content, timestamp, created_at)
VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, ?10, ?11, ?12, ?13)",
rusqlite::params![
id, conv_id, seq as i64, role, content, parts,
tool_call_id, tool_calls, model, status,
reasoning_content, timestamp, created_at
],
)?;
}
}
tx.commit()?;
tracing::info!("v21: 批次 {} 完成({} 对话)", batch_idx, batch.len());
}
conn.execute("INSERT INTO schema_version (version) VALUES (?)", [21])?;
tracing::info!("迁移 v21 完成");
Ok(())
}
```
#### 4.2.2 迁移设计要点
| 要点 | 说明 |
|---|---|
| **幂等安全** | COUNT 探测 + INSERT OR IGNORE;中途崩溃重跑跳过已迁移数据 |
| **分批 commit** | 每 50 对话一批,避免长事务持有 SQLite 写锁导致应用不可用 |
| **迁移期 ID** | `msg_migrated_{conv_id}_{seq}` ——天然唯一(UNIQUE 是 conv_id+seq)、人类可读、零依赖(不调 new_id) |
| **裸 JSON 提取** | 用 `serde_json::Value` 而非 `ChatMessage`,因 df-storage 不依赖 df-ai-core(分层约束) |
| **坏数据跳过** | JSON 解析失败 → warn + continue,不中断迁移 |
| **status 归一化** | None/空 → "active",列语义清晰永不 NULL |
| **created_at 语义** | 有 timestamp 用消息自己的;没有 fallback 到对话 created_at |
#### 4.2.3 迁移耦合点(⚠️ ChatMessage 改名会漏数据)
迁移函数硬编码以下 JSON 字段名,与 `ChatMessage` 的 serde 序列化字段一一对应:
| 硬编码字段名 | ChatMessage 来源 | serde 属性 |
|---|---|---|
| `"role"` | `role: MessageRole` | `#[serde(rename_all = "lowercase")]`(enum) |
| `"content"` | `content: String` | 默认 |
| `"parts"` | `parts: Option<Vec<ContentPart>>` | `skip_serializing_if = "Option::is_none"` |
| `"tool_call_id"` | `tool_call_id: Option<String>` | `skip_serializing_if = "Option::is_none"` |
| `"tool_calls"` | `tool_calls: Option<Vec<ToolCall>>` | `skip_serializing_if = "Option::is_none"` |
| `"model"` | `model: Option<String>` | `default, skip_serializing_if` |
| `"status"` | `status: Option<String>` | `default, skip_serializing_if` |
| `"reasoning_content"` | `reasoning_content: Option<String>` | `default, skip_serializing_if` |
| `"timestamp"` | `timestamp: Option<i64>` | `default, skip_serializing_if` |
> **ChatMessage 如改字段名,必须同步更新迁移函数**。建议在 `types.rs` ChatMessage 定义处加注释标注此耦合点。
#### 4.2.4 双写改造
`save_conversation`(`conversation.rs:135`)改造:
```rust
// 同一事务内:旧列全量写 + 新表增量写
let tx = conn.transaction()?;
// 旧列(回退保险)
tx.execute("UPDATE ai_conversations SET messages = ? WHERE id = ?", params![json, conv_id])?;
// 新表(增量:DELETE 旧 dirty 范围 + INSERT 新)
if let Some((min_seq, max_seq)) = ctx_mgr.dirty_range() {
tx.execute(
"DELETE FROM ai_messages WHERE conversation_id = ? AND seq >= ? AND seq <= ?",
params![conv_id, min_seq, max_seq]
)?;
for (seq, msg) in ctx_mgr.messages()[min_seq..=max_seq].iter().enumerate() {
tx.execute("INSERT INTO ai_messages (...) VALUES (...)", params![...])?;
}
}
tx.commit()?;
```
**风险**:中。双写增加事务开销,但保证一致性。
### Phase 3:切读 + 删旧列
**目标**:读走新表,删旧列,清理双写代码。
| 操作 | 旧实现 | 新实现 |
|---|---|---|
| `restore_from_messages` | 反序列化 messages JSON | `SELECT * FROM ai_messages WHERE conversation_id = ? ORDER BY seq` |
| `clear_messages` | 写空 JSON `[]` | `DELETE FROM ai_messages WHERE conversation_id = ?` |
| `compress_old_messages` | 遍历改 status + 重写 JSON | `UPDATE ai_messages SET status = 'compressed' WHERE conversation_id = ? AND seq < ?` |
| `replace_tool_result_content` | 遍历找 tool_call_id + 重写 JSON | `UPDATE ai_messages SET content = ? WHERE conversation_id = ? AND tool_call_id = ?` |
V22 迁移删 `ai_conversations.messages` 列(SQLite 不支持 DROP COLUMN,需重建表)。
**风险**:中。需全量回归测试对话加载/清空/压缩/编辑。
---
## 五、涉及文件
| 文件 | 改动 |
|---|---|
| `crates/df-storage/src/migrations.rs` | V21 迁移 + V1 建表 SQL 同步补 ai_messages + steps 数组追加 (21, migrate_v21) |
| `crates/df-storage/src/models.rs` | 新增 `AiMessageRecord` struct |
| `crates/df-storage/src/crud/message_repo.rs` | **新建**,AiMessageRepo CRUD(insert_batch / list_by_conversation / delete_range / update_status / update_content_by_tool_call_id) |
| `crates/df-storage/src/crud/mod.rs` | 注册 message_repo 子模块 |
| `src-tauri/src/commands/ai/conversation.rs` | save_conversation 双写改造(Phase 2)/ 切读改造(Phase 3) |
| `src-tauri/src/commands/ai/commands/conversation.rs` | ai_conversation_switch 切读改造(Phase 3) |
| `src-tauri/src/commands/ai/commands/chat.rs` | clear_messages / compress / replace_tool_result_content 改造(Phase 3) |
| `crates/df-ai/src/context.rs` | dirty 标记(Phase 1)+ 适配新存储(Phase 3) |
---
## 六、验收标准
1. **V21 迁移幂等安全**:新库(空表直接建)/老库(全量迁移)/坏数据(JSON 解析失败跳过)均不崩
2. **Phase 2 双写期**:旧列和新表数据一致(可对比校验脚本)
3. **Phase 3 切读后**:对话加载/清空/压缩/编辑行为零回归
4. **长对话(50+ 轮)save 性能显著改善**(增量写入)
5. `cargo check --workspace EXIT 0` + `cargo test` 相关测试通过
---
## 七、风险与对策
| 风险 | 等级 | 对策 |
|---|---|---|
| 迁移函数字段名与 ChatMessage 脱节 | 中 | 迁移耦合点注释 + types.rs 标注 |
| SQLite 不支持 DROP COLUMN | 低 | V22 用重建表方式(CREATE new → INSERT → DROP old → RENAME) |
| 双写事务开销 | 低 | Phase 2 过渡期可接受,Phase 3 清理 |
| ContextManager dirty 标记遗漏 | 中 | 所有变更操作(push/pop/replace/compress/insert)统一设置 dirty |