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

14 KiB

消息拆分存储设计

创建:2026-06-19 | 编号:F-260619-03(消息拆分) | 状态:📐 设计待实施 关联任务:消息拆分存储:ai_messages 表 + 全量迁移 + 三阶段渐进切换(DevFlow 项目) 关联灵感:b2e61a21 消息拆分存储 + 消息级溯源 + 知识脉络网络 关联设计:消息级溯源设计-2026-06-19.md(本设计的前置依赖 + 消费方) 上级索引:../INDEX.md


一、目标

ai_conversations.messages(整个对话 JSON 数组存于单 TEXT 列)拆分为独立的 ai_messages 表(每条消息一行),实现:

  • 增量写入:长对话(50+ 轮)save 不再全量序列化覆盖写
  • O(1) 删除/压缩/单条编辑:clear/compress/replace_tool_result_content 改为 SQL 单行操作
  • 游标分页加载:超长对话按 seq 范围拉取(未来)
  • 持久化消息 ID:为消息级溯源提供 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)——见消息级溯源设计 §三。

2.3 迁移基线

  • 当前最新迁移版本:V19(migrations.rs:45)
  • V20 预留给 F-260619-01(任务关联灵感)
  • 本设计使用 V21(建 ai_messages 表 + 全量迁移)
  • V22 预留给 Phase 3 删旧列

三、数据模型

3.1 新建 ai_messages

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 的序列化开销。

  • ContextManagerdirty_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

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)改造:

// 同一事务内:旧列全量写 + 新表增量写
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