From fc767e1db6fe1efc4810852dd612dbcfae6912b6 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=E7=BB=9D=E5=B0=98?= <237809796@qq.com> Date: Tue, 16 Jun 2026 03:27:01 +0800 Subject: [PATCH] =?UTF-8?q?=E4=BC=98=E5=8C=96:=20R-PD-10=20storage=5Ferr?= =?UTF-8?q?=20DRY=E7=BB=9F=E4=B8=80=20+=20FR-D3=20migrations=20if=E9=93=BE?= =?UTF-8?q?=E8=BD=AC=E6=95=B0=E7=BB=84=E5=BE=AA=E7=8E=AF?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - crud.rs 110处 .map_err(|e| Error::Storage(e.to_string())) 统一 storage_err helper(R-PD-10 DRY) - migrations.rs run() 15个 if current_version(e: E) -> Error { + Error::Storage(e.to_string()) +} + /// 规范化路径用于比较:canonicalize 解析绝对规范路径(失败降级), /// 统一正斜杠 + 小写。与 `df_project::scan::normalize_path` **同算法镜像**(df-storage /// 不依赖 df-project,故独立实现;改动须同步)。防 `C:\a\b` vs `C:/a/b/` 绕过。 @@ -83,11 +89,11 @@ macro_rules! impl_repo { let id = record.id.clone(); let $conn = &*guard; let $rec = &record; - $insert_body.map_err(|e: rusqlite::Error| Error::Storage(e.to_string()))?; + $insert_body.map_err(storage_err)?; Ok(id) }) .await - .map_err(|e| Error::Storage(e.to_string()))? + .map_err(storage_err)? } pub async fn get_by_id(&self, id: &str) -> Result> { @@ -97,15 +103,15 @@ macro_rules! impl_repo { let guard = conn.blocking_lock(); let mut stmt = guard .prepare(&format!("SELECT * FROM {} WHERE id = ?1", $table)) - .map_err(|e| Error::Storage(e.to_string()))?; + .map_err(storage_err)?; let row = stmt .query_row(params![id], |$row| $from_body) .optional() - .map_err(|e| Error::Storage(e.to_string()))?; + .map_err(storage_err)?; Ok(row) }) .await - .map_err(|e| Error::Storage(e.to_string()))? + .map_err(storage_err)? } pub async fn list_all(&self) -> Result> { @@ -114,20 +120,20 @@ macro_rules! impl_repo { let guard = conn.blocking_lock(); let mut stmt = guard .prepare(&format!("SELECT * FROM {} ORDER BY created_at DESC", $table)) - .map_err(|e| Error::Storage(e.to_string()))?; + .map_err(storage_err)?; let rows = stmt .query_map([], |$row| $from_body) - .map_err(|e| Error::Storage(e.to_string()))?; + .map_err(storage_err)?; let mut results = Vec::new(); for r in rows { results.push( - r.map_err(|e| Error::Storage(e.to_string()))?, + r.map_err(storage_err)?, ); } Ok(results) }) .await - .map_err(|e| Error::Storage(e.to_string()))? + .map_err(storage_err)? } pub async fn query(&self, field: &str, value: &str) -> Result> { @@ -140,20 +146,20 @@ macro_rules! impl_repo { let guard = conn.blocking_lock(); let mut stmt = guard .prepare(&sql) - .map_err(|e| Error::Storage(e.to_string()))?; + .map_err(storage_err)?; let rows = stmt .query_map(params![value], |$row| $from_body) - .map_err(|e| Error::Storage(e.to_string()))?; + .map_err(storage_err)?; let mut results = Vec::new(); for r in rows { results.push( - r.map_err(|e| Error::Storage(e.to_string()))?, + r.map_err(storage_err)?, ); } Ok(results) }) .await - .map_err(|e| Error::Storage(e.to_string()))? + .map_err(storage_err)? } pub async fn update_field(&self, id: &str, field: &str, value: &str) -> Result { @@ -167,11 +173,11 @@ macro_rules! impl_repo { let guard = conn.blocking_lock(); let affected = guard .execute(&sql, params![value, now, id]) - .map_err(|e| Error::Storage(e.to_string()))?; + .map_err(storage_err)?; Ok(affected > 0) }) .await - .map_err(|e| Error::Storage(e.to_string()))? + .map_err(storage_err)? } /// 整体更新记录(全可变字段,保留 id 与 created_at),单次原子写 @@ -183,11 +189,11 @@ macro_rules! impl_repo { let $u_conn = &*guard; let $u_rec = &rec; let affected = - $update_body.map_err(|e: rusqlite::Error| Error::Storage(e.to_string()))?; + $update_body.map_err(storage_err)?; Ok(affected > 0) }) .await - .map_err(|e| Error::Storage(e.to_string()))? + .map_err(storage_err)? } pub async fn delete(&self, id: &str) -> Result { @@ -197,11 +203,11 @@ macro_rules! impl_repo { let guard = conn.blocking_lock(); let affected = guard .execute(&format!("DELETE FROM {} WHERE id = ?1", $table), params![id]) - .map_err(|e| Error::Storage(e.to_string()))?; + .map_err(storage_err)?; Ok(affected > 0) }) .await - .map_err(|e| Error::Storage(e.to_string()))? + .map_err(storage_err)? } } }; @@ -237,11 +243,11 @@ impl SettingsRepo { |row| row.get::<_, String>(0), ) .optional() - .map_err(|e| Error::Storage(e.to_string()))?; + .map_err(storage_err)?; Ok(v) }) .await - .map_err(|e| Error::Storage(e.to_string()))? + .map_err(storage_err)? } /// 写 key/value(`INSERT OR REPLACE`),刷新 updated_at @@ -257,11 +263,11 @@ impl SettingsRepo { "INSERT OR REPLACE INTO app_settings (key, value, updated_at) VALUES (?1, ?2, ?3)", params![key, value, now], ) - .map_err(|e| Error::Storage(e.to_string()))?; + .map_err(storage_err)?; Ok(affected > 0) }) .await - .map_err(|e| Error::Storage(e.to_string()))? + .map_err(storage_err)? } /// 取全部 key/value @@ -271,18 +277,18 @@ impl SettingsRepo { let guard = conn.blocking_lock(); let mut stmt = guard .prepare("SELECT key, value FROM app_settings") - .map_err(|e| Error::Storage(e.to_string()))?; + .map_err(storage_err)?; let rows = stmt .query_map([], |row| Ok((row.get::<_, String>(0)?, row.get::<_, String>(1)?))) - .map_err(|e| Error::Storage(e.to_string()))?; + .map_err(storage_err)?; let mut results = Vec::new(); for r in rows { - results.push(r.map_err(|e| Error::Storage(e.to_string()))?); + results.push(r.map_err(storage_err)?); } Ok(results) }) .await - .map_err(|e| Error::Storage(e.to_string()))? + .map_err(storage_err)? } /// 删除 key,返回是否实际删除 @@ -293,11 +299,11 @@ impl SettingsRepo { let guard = conn.blocking_lock(); let affected = guard .execute("DELETE FROM app_settings WHERE key = ?1", params![key]) - .map_err(|e| Error::Storage(e.to_string()))?; + .map_err(storage_err)?; Ok(affected > 0) }) .await - .map_err(|e| Error::Storage(e.to_string()))? + .map_err(storage_err)? } } @@ -606,18 +612,18 @@ impl ProjectRepo { let guard = conn.blocking_lock(); let mut stmt = guard .prepare("SELECT id, name, description, status, idea_id, path, stack, created_at, updated_at FROM projects WHERE deleted_at IS NULL ORDER BY created_at DESC") - .map_err(|e| Error::Storage(e.to_string()))?; + .map_err(storage_err)?; let rows = stmt .query_map([], |row| project_from_row(row)) - .map_err(|e| Error::Storage(e.to_string()))?; + .map_err(storage_err)?; let mut results = Vec::new(); for r in rows { - results.push(r.map_err(|e| Error::Storage(e.to_string()))?); + results.push(r.map_err(storage_err)?); } Ok(results) }) .await - .map_err(|e| Error::Storage(e.to_string()))? + .map_err(storage_err)? } /// 软删:标记 deleted_at(进回收站,可恢复)。仅作用于未删项目,返回是否命中。 @@ -632,11 +638,11 @@ impl ProjectRepo { "UPDATE projects SET deleted_at = ?1, updated_at = ?1 WHERE id = ?2 AND deleted_at IS NULL", params![now, id], ) - .map_err(|e| Error::Storage(e.to_string()))?; + .map_err(storage_err)?; Ok(affected > 0) }) .await - .map_err(|e| Error::Storage(e.to_string()))? + .map_err(storage_err)? } /// 恢复:清 deleted_at(从回收站还原) @@ -651,11 +657,11 @@ impl ProjectRepo { "UPDATE projects SET deleted_at = NULL, updated_at = ?1 WHERE id = ?2 AND deleted_at IS NOT NULL", params![now, id], ) - .map_err(|e| Error::Storage(e.to_string()))?; + .map_err(storage_err)?; Ok(affected > 0) }) .await - .map_err(|e| Error::Storage(e.to_string()))? + .map_err(storage_err)? } /// 查找已绑定该规范化路径的项目(排除 exclude_id 自身)。无冲突返回 None。 @@ -686,18 +692,18 @@ impl ProjectRepo { let guard = conn.blocking_lock(); let mut stmt = guard .prepare("SELECT id, name, description, status, idea_id, path, stack, created_at, updated_at FROM projects WHERE deleted_at IS NOT NULL ORDER BY updated_at DESC") - .map_err(|e| Error::Storage(e.to_string()))?; + .map_err(storage_err)?; let rows = stmt .query_map([], |row| project_from_row(row)) - .map_err(|e| Error::Storage(e.to_string()))?; + .map_err(storage_err)?; let mut results = Vec::new(); for r in rows { - results.push(r.map_err(|e| Error::Storage(e.to_string()))?); + results.push(r.map_err(storage_err)?); } Ok(results) }) .await - .map_err(|e| Error::Storage(e.to_string()))? + .map_err(storage_err)? } /// 彻底删除:事务级联删 branches→releases→tasks→projects(不可恢复) @@ -709,21 +715,21 @@ impl ProjectRepo { let id = id.to_owned(); tokio::task::spawn_blocking(move || { let mut guard = conn.blocking_lock(); - let tx = guard.transaction().map_err(|e| Error::Storage(e.to_string()))?; + let tx = guard.transaction().map_err(storage_err)?; tx.execute("DELETE FROM branches WHERE project_id = ?1", params![id]) - .map_err(|e| Error::Storage(e.to_string()))?; + .map_err(storage_err)?; tx.execute("DELETE FROM releases WHERE project_id = ?1", params![id]) - .map_err(|e| Error::Storage(e.to_string()))?; + .map_err(storage_err)?; tx.execute("DELETE FROM tasks WHERE project_id = ?1", params![id]) - .map_err(|e| Error::Storage(e.to_string()))?; + .map_err(storage_err)?; let affected = tx .execute("DELETE FROM projects WHERE id = ?1", params![id]) - .map_err(|e| Error::Storage(e.to_string()))?; - tx.commit().map_err(|e| Error::Storage(e.to_string()))?; + .map_err(storage_err)?; + tx.commit().map_err(storage_err)?; Ok(affected > 0) }) .await - .map_err(|e| Error::Storage(e.to_string()))? + .map_err(storage_err)? } } @@ -767,18 +773,18 @@ impl TaskRepo { let guard = conn.blocking_lock(); let mut stmt = guard .prepare("SELECT id, project_id, title, description, status, priority, branch_name, assignee, workflow_def_id, base_branch, review_rounds, created_at, updated_at FROM tasks WHERE deleted_at IS NULL ORDER BY created_at DESC") - .map_err(|e| Error::Storage(e.to_string()))?; + .map_err(storage_err)?; let rows = stmt .query_map([], |row| task_from_row(row)) - .map_err(|e| Error::Storage(e.to_string()))?; + .map_err(storage_err)?; let mut results = Vec::new(); for r in rows { - results.push(r.map_err(|e| Error::Storage(e.to_string()))?); + results.push(r.map_err(storage_err)?); } Ok(results) }) .await - .map_err(|e| Error::Storage(e.to_string()))? + .map_err(storage_err)? } /// 软删:标记 deleted_at(进回收站,可恢复)。仅作用于未删任务,返回是否命中。 @@ -794,11 +800,11 @@ impl TaskRepo { "UPDATE tasks SET deleted_at = ?1, updated_at = ?1 WHERE id = ?2 AND deleted_at IS NULL", params![now, id], ) - .map_err(|e| Error::Storage(e.to_string()))?; + .map_err(storage_err)?; Ok(affected > 0) }) .await - .map_err(|e| Error::Storage(e.to_string()))? + .map_err(storage_err)? } /// 恢复:清 deleted_at(从回收站还原)。仅作用于已删任务,返回是否命中。 @@ -814,11 +820,11 @@ impl TaskRepo { "UPDATE tasks SET deleted_at = NULL, updated_at = ?1 WHERE id = ?2 AND deleted_at IS NOT NULL", params![now, id], ) - .map_err(|e| Error::Storage(e.to_string()))?; + .map_err(storage_err)?; Ok(affected > 0) }) .await - .map_err(|e| Error::Storage(e.to_string()))? + .map_err(storage_err)? } /// 原子推进任务状态(任务推进链 F-260616-02 唯一 status 写入路径) @@ -861,22 +867,22 @@ impl TaskRepo { }; let affected = guard .execute(sql, params![new_status, now, id, expected]) - .map_err(|e| Error::Storage(e.to_string()))?; + .map_err(storage_err)?; if affected == 0 { return Ok(None); } // 回读更新后的记录(含新 status / 累加后的 review_rounds / 新 updated_at)。 let mut stmt = guard .prepare("SELECT id, project_id, title, description, status, priority, branch_name, assignee, workflow_def_id, base_branch, review_rounds, created_at, updated_at FROM tasks WHERE id = ?1") - .map_err(|e| Error::Storage(e.to_string()))?; + .map_err(storage_err)?; let row = stmt .query_row(params![id], |row| task_from_row(row)) .optional() - .map_err(|e| Error::Storage(e.to_string()))?; + .map_err(storage_err)?; Ok(row) }) .await - .map_err(|e| Error::Storage(e.to_string()))? + .map_err(storage_err)? } /// 列出回收站(deleted_at IS NOT NULL),按更新时间(≈删除时间)降序。对标 ProjectRepo::list_deleted。 @@ -889,18 +895,18 @@ impl TaskRepo { let guard = conn.blocking_lock(); let mut stmt = guard .prepare("SELECT id, project_id, title, description, status, priority, branch_name, assignee, workflow_def_id, base_branch, review_rounds, created_at, updated_at FROM tasks WHERE deleted_at IS NOT NULL ORDER BY updated_at DESC") - .map_err(|e| Error::Storage(e.to_string()))?; + .map_err(storage_err)?; let rows = stmt .query_map([], |row| task_from_row(row)) - .map_err(|e| Error::Storage(e.to_string()))?; + .map_err(storage_err)?; let mut results = Vec::new(); for r in rows { - results.push(r.map_err(|e| Error::Storage(e.to_string()))?); + results.push(r.map_err(storage_err)?); } Ok(results) }) .await - .map_err(|e| Error::Storage(e.to_string()))? + .map_err(storage_err)? } } @@ -1198,15 +1204,15 @@ impl AiToolExecutionRepo { .prepare( "SELECT * FROM ai_tool_executions WHERE tool_call_id = ?1 ORDER BY requested_at DESC LIMIT 1", ) - .map_err(|e| Error::Storage(e.to_string()))?; + .map_err(storage_err)?; let row = stmt .query_row(params![tid], |row| ai_tool_execution_from_row(row)) .optional() - .map_err(|e| Error::Storage(e.to_string()))?; + .map_err(storage_err)?; Ok(row) }) .await - .map_err(|e| Error::Storage(e.to_string()))? + .map_err(storage_err)? } /// 列出所有 status=pending 的审计行(启动重建 pending_approvals 用) @@ -1218,18 +1224,18 @@ impl AiToolExecutionRepo { let guard = conn.blocking_lock(); let mut stmt = guard .prepare("SELECT * FROM ai_tool_executions WHERE status = 'pending' ORDER BY requested_at ASC") - .map_err(|e| Error::Storage(e.to_string()))?; + .map_err(storage_err)?; let rows = stmt .query_map([], |row| ai_tool_execution_from_row(row)) - .map_err(|e| Error::Storage(e.to_string()))?; + .map_err(storage_err)?; let mut results = Vec::new(); for r in rows { - results.push(r.map_err(|e| Error::Storage(e.to_string()))?); + results.push(r.map_err(storage_err)?); } Ok(results) }) .await - .map_err(|e| Error::Storage(e.to_string()))? + .map_err(storage_err)? } } @@ -1281,28 +1287,28 @@ impl KnowledgeRepo { if let Some(k) = &kind { let mut stmt = guard .prepare(&format!("SELECT {KNOWLEDGE_COLS} FROM knowledges WHERE status = 'published' AND (title LIKE ?1 OR content LIKE ?2) AND kind = ?3 ORDER BY reuse_count DESC LIMIT ?4")) - .map_err(|e| Error::Storage(e.to_string()))?; + .map_err(storage_err)?; let rows = stmt .query_map(params![pattern, pattern, k, limit_i], |row| knowledge_from_row(row)) - .map_err(|e| Error::Storage(e.to_string()))?; + .map_err(storage_err)?; for r in rows { - results.push(r.map_err(|e| Error::Storage(e.to_string()))?); + results.push(r.map_err(storage_err)?); } } else { let mut stmt = guard .prepare(&format!("SELECT {KNOWLEDGE_COLS} FROM knowledges WHERE status = 'published' AND (title LIKE ?1 OR content LIKE ?2) ORDER BY reuse_count DESC LIMIT ?3")) - .map_err(|e| Error::Storage(e.to_string()))?; + .map_err(storage_err)?; let rows = stmt .query_map(params![pattern, pattern, limit_i], |row| knowledge_from_row(row)) - .map_err(|e| Error::Storage(e.to_string()))?; + .map_err(storage_err)?; for r in rows { - results.push(r.map_err(|e| Error::Storage(e.to_string()))?); + results.push(r.map_err(storage_err)?); } } Ok(results) }) .await - .map_err(|e| Error::Storage(e.to_string()))? + .map_err(storage_err)? } /// 按状态列出(审核收件箱用): 按 confidence 语义排序(high>medium>low),次按 created_at @@ -1325,18 +1331,18 @@ impl KnowledgeRepo { ELSE 0 END DESC, created_at DESC", ) - .map_err(|e| Error::Storage(e.to_string()))?; + .map_err(storage_err)?; let rows = stmt .query_map(params![status], |row| knowledge_from_row(row)) - .map_err(|e| Error::Storage(e.to_string()))?; + .map_err(storage_err)?; let mut results = Vec::new(); for r in rows { - results.push(r.map_err(|e| Error::Storage(e.to_string()))?); + results.push(r.map_err(storage_err)?); } Ok(results) }) .await - .map_err(|e| Error::Storage(e.to_string()))? + .map_err(storage_err)? } /// 复用计数 +1(SQL 行级原子操作,并发安全) @@ -1351,11 +1357,11 @@ impl KnowledgeRepo { "UPDATE knowledges SET reuse_count = reuse_count + 1, updated_at = ?1 WHERE id = ?2", params![now, id], ) - .map_err(|e| Error::Storage(e.to_string()))?; + .map_err(storage_err)?; Ok(affected > 0) }) .await - .map_err(|e| Error::Storage(e.to_string()))? + .map_err(storage_err)? } /// 写入向量嵌入(BLOB = Vec 小端字节序列化) @@ -1372,11 +1378,11 @@ impl KnowledgeRepo { "UPDATE knowledges SET embedding = ?1 WHERE id = ?2", params![blob, id], ) - .map_err(|e| Error::Storage(e.to_string()))?; + .map_err(storage_err)?; Ok(affected > 0) }) .await - .map_err(|e| Error::Storage(e.to_string()))? + .map_err(storage_err)? } /// 向量检索: 加载全部 published 且有 embedding 的记录,纯 Rust 余弦相似度取 top-N @@ -1404,18 +1410,18 @@ impl KnowledgeRepo { .prepare(&format!( "SELECT {KNOWLEDGE_COLS_WITH_EMBEDDING} FROM knowledges WHERE status = 'published' AND embedding IS NOT NULL" )) - .map_err(|e| Error::Storage(e.to_string()))?; + .map_err(storage_err)?; let rows = stmt .query_map([], |row| { let rec = knowledge_from_row(row)?; let blob: Vec = row.get("embedding")?; Ok((rec, blob)) }) - .map_err(|e| Error::Storage(e.to_string()))?; + .map_err(storage_err)?; let mut scored: Vec<(KnowledgeRecord, f32)> = Vec::new(); for r in rows { - let (rec, blob) = r.map_err(|e| Error::Storage(e.to_string()))?; + let (rec, blob) = r.map_err(storage_err)?; let emb = blob_to_f32s(&blob); // 维度不匹配(换过 embedding 模型的旧向量)跳过 if emb.len() != query_vec.len() { @@ -1429,7 +1435,7 @@ impl KnowledgeRepo { Ok(scored) }) .await - .map_err(|e| Error::Storage(e.to_string()))? + .map_err(storage_err)? } /// 列出非归档知识(全部 status != 'archived'),按 confidence 语义排序 @@ -1449,18 +1455,18 @@ impl KnowledgeRepo { ELSE 0 END DESC, created_at DESC", ) - .map_err(|e| Error::Storage(e.to_string()))?; + .map_err(storage_err)?; let rows = stmt .query_map([], |row| knowledge_from_row(row)) - .map_err(|e| Error::Storage(e.to_string()))?; + .map_err(storage_err)?; let mut results = Vec::new(); for r in rows { - results.push(r.map_err(|e| Error::Storage(e.to_string()))?); + results.push(r.map_err(storage_err)?); } Ok(results) }) .await - .map_err(|e| Error::Storage(e.to_string()))? + .map_err(storage_err)? } /// 热门知识(已发布,按复用次数降序) @@ -1471,18 +1477,18 @@ impl KnowledgeRepo { let guard = conn.blocking_lock(); let mut stmt = guard .prepare("SELECT id,kind,title,content,tags,status,confidence,reuse_count,verified,source_project,source_ref,reasoning,created_at,updated_at FROM knowledges WHERE status = 'published' ORDER BY reuse_count DESC LIMIT ?1") - .map_err(|e| Error::Storage(e.to_string()))?; + .map_err(storage_err)?; let rows = stmt .query_map(params![limit_i], |row| knowledge_from_row(row)) - .map_err(|e| Error::Storage(e.to_string()))?; + .map_err(storage_err)?; let mut results = Vec::new(); for r in rows { - results.push(r.map_err(|e| Error::Storage(e.to_string()))?); + results.push(r.map_err(storage_err)?); } Ok(results) }) .await - .map_err(|e| Error::Storage(e.to_string()))? + .map_err(storage_err)? } } @@ -1516,18 +1522,18 @@ impl KnowledgeEventsRepo { let guard = conn.blocking_lock(); let mut stmt = guard .prepare("SELECT id,knowledge_id,event_type,source_ref,context_json,timestamp FROM knowledge_events WHERE knowledge_id = ?1 ORDER BY timestamp ASC") - .map_err(|e| Error::Storage(e.to_string()))?; + .map_err(storage_err)?; let rows = stmt .query_map(params![knowledge_id], |row| knowledge_event_from_row(row)) - .map_err(|e| Error::Storage(e.to_string()))?; + .map_err(storage_err)?; let mut results = Vec::new(); for r in rows { - results.push(r.map_err(|e| Error::Storage(e.to_string()))?); + results.push(r.map_err(storage_err)?); } Ok(results) }) .await - .map_err(|e| Error::Storage(e.to_string()))? + .map_err(storage_err)? } /// 按知识 ID + event_type 查询最近 N 条(如引用记录翻页: type=referenced) @@ -1545,18 +1551,18 @@ impl KnowledgeEventsRepo { let guard = conn.blocking_lock(); let mut stmt = guard .prepare("SELECT id,knowledge_id,event_type,source_ref,context_json,timestamp FROM knowledge_events WHERE knowledge_id = ?1 AND event_type = ?2 ORDER BY timestamp DESC LIMIT ?3") - .map_err(|e| Error::Storage(e.to_string()))?; + .map_err(storage_err)?; let rows = stmt .query_map(params![knowledge_id, event_type, limit_i], |row| knowledge_event_from_row(row)) - .map_err(|e| Error::Storage(e.to_string()))?; + .map_err(storage_err)?; let mut results = Vec::new(); for r in rows { - results.push(r.map_err(|e| Error::Storage(e.to_string()))?); + results.push(r.map_err(storage_err)?); } Ok(results) }) .await - .map_err(|e| Error::Storage(e.to_string()))? + .map_err(storage_err)? } } @@ -1578,11 +1584,11 @@ impl AiConversationRepo { "UPDATE ai_conversations SET messages = '[]', prompt_tokens = 0, completion_tokens = 0, updated_at = ?1 WHERE id = ?2", params![now, id], ) - .map_err(|e| Error::Storage(e.to_string()))?; + .map_err(storage_err)?; Ok(affected > 0) }) .await - .map_err(|e| Error::Storage(e.to_string()))? + .map_err(storage_err)? } /// 设置归档标记(仅改 archived,不动 updated_at) @@ -1599,11 +1605,11 @@ impl AiConversationRepo { "UPDATE ai_conversations SET archived = ?1 WHERE id = ?2", params![if archived { 1 } else { 0 }, id], ) - .map_err(|e| Error::Storage(e.to_string()))?; + .map_err(storage_err)?; Ok(affected > 0) }) .await - .map_err(|e| Error::Storage(e.to_string()))? + .map_err(storage_err)? } } diff --git a/crates/df-storage/src/migrations.rs b/crates/df-storage/src/migrations.rs index d6606e9..4310f43 100644 --- a/crates/df-storage/src/migrations.rs +++ b/crates/df-storage/src/migrations.rs @@ -32,64 +32,30 @@ pub fn run(conn: &Connection) -> Result<()> { ) .unwrap_or(0); - if current_version < 1 { - migrate_v1(conn)?; - } + // 迁移步骤链: 顺序执行,跳过已应用的版本(current_version < N 才跑)。 + // 新增版本时,在此数组追加一项 (N, migrate_vN) 即可,无需改逻辑。 + let steps: [(i32, fn(&Connection) -> Result<()>); 15] = [ + (1, migrate_v1), + (2, migrate_v2), + (3, migrate_v3), + (4, migrate_v4), + (5, migrate_v5), + (6, migrate_v6), + (7, migrate_v7), + (8, migrate_v8), + (9, migrate_v9), + (10, migrate_v10), + (11, migrate_v11), + (12, migrate_v12), + (13, migrate_v13), + (14, migrate_v14), + (15, migrate_v15), + ]; - if current_version < 2 { - migrate_v2(conn)?; - } - - if current_version < 3 { - migrate_v3(conn)?; - } - - if current_version < 4 { - migrate_v4(conn)?; - } - - if current_version < 5 { - migrate_v5(conn)?; - } - - if current_version < 6 { - migrate_v6(conn)?; - } - - if current_version < 7 { - migrate_v7(conn)?; - } - - if current_version < 8 { - migrate_v8(conn)?; - } - - if current_version < 9 { - migrate_v9(conn)?; - } - - if current_version < 10 { - migrate_v10(conn)?; - } - - if current_version < 11 { - migrate_v11(conn)?; - } - - if current_version < 12 { - migrate_v12(conn)?; - } - - if current_version < 13 { - migrate_v13(conn)?; - } - - if current_version < 14 { - migrate_v14(conn)?; - } - - if current_version < 15 { - migrate_v15(conn)?; + for (version, migrate_fn) in steps { + if current_version < version { + migrate_fn(conn)?; + } } Ok(())