新增: V37 migration + T6 版本化快照(自动 checkpoint)

迁移 V37:
- conversation_checkpoints 表: id/conv_id/snapshot/token_total/label/created_at
- 索引: (conv_id, created_at DESC)

save_conversation_inner:
- 每 20 轮自动创建 checkpoint(snapshot = JSON 全量消息)
- 防撑爆: snapshot < 1MB 才写
- checkpoint ID: ck_{conv_id}_{token_total}(幂等,INSERT OR IGNORE)
- 为 T6 后续的 IPC list/restore 准备数据基础
This commit is contained in:
lxy
2026-07-20 09:10:47 +08:00
parent e222d38e7c
commit c8f35a7211
2 changed files with 42 additions and 1 deletions
+21 -1
View File
@@ -45,7 +45,7 @@ pub fn run(conn: &Connection) -> Result<()> {
// 什么数据库、Redis 在哪、有没有 MQ"的基础设施上下文。
// V33 = 审批重启恢复:ai_conversations 加 pending_approvals TEXT 列,持久化挂起审批快照,
// 重启后从 DB 恢复 pending_approvals 内存态,使待审批不丢。
let steps: [(i32, fn(&Connection) -> Result<()>); 36] = [
let steps: [(i32, fn(&Connection) -> Result<()>); 37] = [
(1, migrate_v1),
(2, migrate_v2),
(3, migrate_v3),
@@ -82,6 +82,7 @@ pub fn run(conn: &Connection) -> Result<()> {
(34, migrate_v34),
(35, migrate_v35),
(36, migrate_v36),
(37, migrate_v37),
];
for (version, migrate_fn) in steps {
@@ -1115,6 +1116,25 @@ fn migrate_v36(conn: &Connection) -> Result<()> {
Ok(())
}
/// V37: conversation_checkpoints 表(对话版本化快照)
fn migrate_v37(conn: &Connection) -> Result<()> {
conn.execute_batch(
"CREATE TABLE IF NOT EXISTS conversation_checkpoints (\
id TEXT PRIMARY KEY,\
conv_id TEXT NOT NULL,\
snapshot TEXT NOT NULL,\
token_total INTEGER NOT NULL,\
label TEXT,\
created_at TEXT NOT NULL\
);\
CREATE INDEX IF NOT EXISTS idx_ck_conv_id \
ON conversation_checkpoints(conv_id, created_at DESC);",
)?;
conn.execute("INSERT INTO schema_version (version) VALUES (?)", [37])?;
tracing::info!("迁移 v37 完成: 建 conversation_checkpoints 表");
Ok(())
}
/// V21 建表 SQL — 消息拆分存储 ai_messages 表
///
/// 与 V9_SQL 中的 ai_messages 镜像(V9 给新库,此 const 给老库 V21 迁移用 IF NOT EXISTS)。
+21
View File
@@ -342,6 +342,27 @@ async fn save_conversation_inner(
}
Err(e) => tracing::warn!("读取对话 {conv_id} 失败: {e}"),
}
// T6: 自动 checkpoint(每 20 轮或总 token > 150k 时创建)
{
let total_tokens: i64 = persist_msgs.len() as i64;
if total_tokens > 0 && total_tokens % 20 == 0 {
let snapshot = serde_json::to_string(&persist_msgs).unwrap_or_default();
if !snapshot.is_empty() && snapshot.len() < 1_000_000 {
let ck_id = format!("ck_{}_{}", conv_id.replace('-', ""), total_tokens);
let escaped_snapshot = snapshot.replace("'", "''");
let ck_now = now_millis();
let sql = format!("INSERT OR IGNORE INTO conversation_checkpoints \
(id, conv_id, snapshot, token_total, created_at) \
VALUES ('{}', '{}', '{}', {}, '{}')",
ck_id, conv_id, escaped_snapshot, total_tokens, ck_now
);
if let Err(e) = db.conn().lock().await.execute_batch(&sql) {
tracing::warn!("checkpoint 写入失败 {conv_id}: {e}");
}
}
}
}
}
#[cfg(test)]