优化: R-PD-10 storage_err DRY统一 + FR-D3 migrations if链转数组循环

- 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<N 块改 step 数组循环,零行为变更(FR-D3),新增版本追加一行即可
This commit is contained in:
2026-06-16 03:27:01 +08:00
parent 49e61d4153
commit fc767e1db6
2 changed files with 139 additions and 167 deletions

View File

@@ -16,6 +16,12 @@ use crate::models::{
TaskRecord, WorkflowRecord, TaskRecord, WorkflowRecord,
}; };
/// rusqlite 错误 → df-core `Error::Storage` 统一包装(R-PD-10 DRY,
/// 替代散落的 `.map_err(storage_err)`)。
fn storage_err<E: std::string::ToString>(e: E) -> Error {
Error::Storage(e.to_string())
}
/// 规范化路径用于比较:canonicalize 解析绝对规范路径(失败降级), /// 规范化路径用于比较:canonicalize 解析绝对规范路径(失败降级),
/// 统一正斜杠 + 小写。与 `df_project::scan::normalize_path` **同算法镜像**(df-storage /// 统一正斜杠 + 小写。与 `df_project::scan::normalize_path` **同算法镜像**(df-storage
/// 不依赖 df-project,故独立实现;改动须同步)。防 `C:\a\b` vs `C:/a/b/` 绕过。 /// 不依赖 df-project,故独立实现;改动须同步)。防 `C:\a\b` vs `C:/a/b/` 绕过。
@@ -83,11 +89,11 @@ macro_rules! impl_repo {
let id = record.id.clone(); let id = record.id.clone();
let $conn = &*guard; let $conn = &*guard;
let $rec = &record; let $rec = &record;
$insert_body.map_err(|e: rusqlite::Error| Error::Storage(e.to_string()))?; $insert_body.map_err(storage_err)?;
Ok(id) Ok(id)
}) })
.await .await
.map_err(|e| Error::Storage(e.to_string()))? .map_err(storage_err)?
} }
pub async fn get_by_id(&self, id: &str) -> Result<Option<$record>> { pub async fn get_by_id(&self, id: &str) -> Result<Option<$record>> {
@@ -97,15 +103,15 @@ macro_rules! impl_repo {
let guard = conn.blocking_lock(); let guard = conn.blocking_lock();
let mut stmt = guard let mut stmt = guard
.prepare(&format!("SELECT * FROM {} WHERE id = ?1", $table)) .prepare(&format!("SELECT * FROM {} WHERE id = ?1", $table))
.map_err(|e| Error::Storage(e.to_string()))?; .map_err(storage_err)?;
let row = stmt let row = stmt
.query_row(params![id], |$row| $from_body) .query_row(params![id], |$row| $from_body)
.optional() .optional()
.map_err(|e| Error::Storage(e.to_string()))?; .map_err(storage_err)?;
Ok(row) Ok(row)
}) })
.await .await
.map_err(|e| Error::Storage(e.to_string()))? .map_err(storage_err)?
} }
pub async fn list_all(&self) -> Result<Vec<$record>> { pub async fn list_all(&self) -> Result<Vec<$record>> {
@@ -114,20 +120,20 @@ macro_rules! impl_repo {
let guard = conn.blocking_lock(); let guard = conn.blocking_lock();
let mut stmt = guard let mut stmt = guard
.prepare(&format!("SELECT * FROM {} ORDER BY created_at DESC", $table)) .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 let rows = stmt
.query_map([], |$row| $from_body) .query_map([], |$row| $from_body)
.map_err(|e| Error::Storage(e.to_string()))?; .map_err(storage_err)?;
let mut results = Vec::new(); let mut results = Vec::new();
for r in rows { for r in rows {
results.push( results.push(
r.map_err(|e| Error::Storage(e.to_string()))?, r.map_err(storage_err)?,
); );
} }
Ok(results) Ok(results)
}) })
.await .await
.map_err(|e| Error::Storage(e.to_string()))? .map_err(storage_err)?
} }
pub async fn query(&self, field: &str, value: &str) -> Result<Vec<$record>> { pub async fn query(&self, field: &str, value: &str) -> Result<Vec<$record>> {
@@ -140,20 +146,20 @@ macro_rules! impl_repo {
let guard = conn.blocking_lock(); let guard = conn.blocking_lock();
let mut stmt = guard let mut stmt = guard
.prepare(&sql) .prepare(&sql)
.map_err(|e| Error::Storage(e.to_string()))?; .map_err(storage_err)?;
let rows = stmt let rows = stmt
.query_map(params![value], |$row| $from_body) .query_map(params![value], |$row| $from_body)
.map_err(|e| Error::Storage(e.to_string()))?; .map_err(storage_err)?;
let mut results = Vec::new(); let mut results = Vec::new();
for r in rows { for r in rows {
results.push( results.push(
r.map_err(|e| Error::Storage(e.to_string()))?, r.map_err(storage_err)?,
); );
} }
Ok(results) Ok(results)
}) })
.await .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<bool> { pub async fn update_field(&self, id: &str, field: &str, value: &str) -> Result<bool> {
@@ -167,11 +173,11 @@ macro_rules! impl_repo {
let guard = conn.blocking_lock(); let guard = conn.blocking_lock();
let affected = guard let affected = guard
.execute(&sql, params![value, now, id]) .execute(&sql, params![value, now, id])
.map_err(|e| Error::Storage(e.to_string()))?; .map_err(storage_err)?;
Ok(affected > 0) Ok(affected > 0)
}) })
.await .await
.map_err(|e| Error::Storage(e.to_string()))? .map_err(storage_err)?
} }
/// 整体更新记录(全可变字段,保留 id 与 created_at单次原子写 /// 整体更新记录(全可变字段,保留 id 与 created_at单次原子写
@@ -183,11 +189,11 @@ macro_rules! impl_repo {
let $u_conn = &*guard; let $u_conn = &*guard;
let $u_rec = &rec; let $u_rec = &rec;
let affected = let affected =
$update_body.map_err(|e: rusqlite::Error| Error::Storage(e.to_string()))?; $update_body.map_err(storage_err)?;
Ok(affected > 0) Ok(affected > 0)
}) })
.await .await
.map_err(|e| Error::Storage(e.to_string()))? .map_err(storage_err)?
} }
pub async fn delete(&self, id: &str) -> Result<bool> { pub async fn delete(&self, id: &str) -> Result<bool> {
@@ -197,11 +203,11 @@ macro_rules! impl_repo {
let guard = conn.blocking_lock(); let guard = conn.blocking_lock();
let affected = guard let affected = guard
.execute(&format!("DELETE FROM {} WHERE id = ?1", $table), params![id]) .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) Ok(affected > 0)
}) })
.await .await
.map_err(|e| Error::Storage(e.to_string()))? .map_err(storage_err)?
} }
} }
}; };
@@ -237,11 +243,11 @@ impl SettingsRepo {
|row| row.get::<_, String>(0), |row| row.get::<_, String>(0),
) )
.optional() .optional()
.map_err(|e| Error::Storage(e.to_string()))?; .map_err(storage_err)?;
Ok(v) Ok(v)
}) })
.await .await
.map_err(|e| Error::Storage(e.to_string()))? .map_err(storage_err)?
} }
/// 写 key/value(`INSERT OR REPLACE`),刷新 updated_at /// 写 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)", "INSERT OR REPLACE INTO app_settings (key, value, updated_at) VALUES (?1, ?2, ?3)",
params![key, value, now], params![key, value, now],
) )
.map_err(|e| Error::Storage(e.to_string()))?; .map_err(storage_err)?;
Ok(affected > 0) Ok(affected > 0)
}) })
.await .await
.map_err(|e| Error::Storage(e.to_string()))? .map_err(storage_err)?
} }
/// 取全部 key/value /// 取全部 key/value
@@ -271,18 +277,18 @@ impl SettingsRepo {
let guard = conn.blocking_lock(); let guard = conn.blocking_lock();
let mut stmt = guard let mut stmt = guard
.prepare("SELECT key, value FROM app_settings") .prepare("SELECT key, value FROM app_settings")
.map_err(|e| Error::Storage(e.to_string()))?; .map_err(storage_err)?;
let rows = stmt let rows = stmt
.query_map([], |row| Ok((row.get::<_, String>(0)?, row.get::<_, String>(1)?))) .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(); let mut results = Vec::new();
for r in rows { for r in rows {
results.push(r.map_err(|e| Error::Storage(e.to_string()))?); results.push(r.map_err(storage_err)?);
} }
Ok(results) Ok(results)
}) })
.await .await
.map_err(|e| Error::Storage(e.to_string()))? .map_err(storage_err)?
} }
/// 删除 key,返回是否实际删除 /// 删除 key,返回是否实际删除
@@ -293,11 +299,11 @@ impl SettingsRepo {
let guard = conn.blocking_lock(); let guard = conn.blocking_lock();
let affected = guard let affected = guard
.execute("DELETE FROM app_settings WHERE key = ?1", params![key]) .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) Ok(affected > 0)
}) })
.await .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 guard = conn.blocking_lock();
let mut stmt = guard 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") .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 let rows = stmt
.query_map([], |row| project_from_row(row)) .query_map([], |row| project_from_row(row))
.map_err(|e| Error::Storage(e.to_string()))?; .map_err(storage_err)?;
let mut results = Vec::new(); let mut results = Vec::new();
for r in rows { for r in rows {
results.push(r.map_err(|e| Error::Storage(e.to_string()))?); results.push(r.map_err(storage_err)?);
} }
Ok(results) Ok(results)
}) })
.await .await
.map_err(|e| Error::Storage(e.to_string()))? .map_err(storage_err)?
} }
/// 软删:标记 deleted_at(进回收站,可恢复)。仅作用于未删项目,返回是否命中。 /// 软删:标记 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", "UPDATE projects SET deleted_at = ?1, updated_at = ?1 WHERE id = ?2 AND deleted_at IS NULL",
params![now, id], params![now, id],
) )
.map_err(|e| Error::Storage(e.to_string()))?; .map_err(storage_err)?;
Ok(affected > 0) Ok(affected > 0)
}) })
.await .await
.map_err(|e| Error::Storage(e.to_string()))? .map_err(storage_err)?
} }
/// 恢复:清 deleted_at(从回收站还原) /// 恢复:清 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", "UPDATE projects SET deleted_at = NULL, updated_at = ?1 WHERE id = ?2 AND deleted_at IS NOT NULL",
params![now, id], params![now, id],
) )
.map_err(|e| Error::Storage(e.to_string()))?; .map_err(storage_err)?;
Ok(affected > 0) Ok(affected > 0)
}) })
.await .await
.map_err(|e| Error::Storage(e.to_string()))? .map_err(storage_err)?
} }
/// 查找已绑定该规范化路径的项目(排除 exclude_id 自身)。无冲突返回 None。 /// 查找已绑定该规范化路径的项目(排除 exclude_id 自身)。无冲突返回 None。
@@ -686,18 +692,18 @@ impl ProjectRepo {
let guard = conn.blocking_lock(); let guard = conn.blocking_lock();
let mut stmt = guard 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") .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 let rows = stmt
.query_map([], |row| project_from_row(row)) .query_map([], |row| project_from_row(row))
.map_err(|e| Error::Storage(e.to_string()))?; .map_err(storage_err)?;
let mut results = Vec::new(); let mut results = Vec::new();
for r in rows { for r in rows {
results.push(r.map_err(|e| Error::Storage(e.to_string()))?); results.push(r.map_err(storage_err)?);
} }
Ok(results) Ok(results)
}) })
.await .await
.map_err(|e| Error::Storage(e.to_string()))? .map_err(storage_err)?
} }
/// 彻底删除:事务级联删 branches→releases→tasks→projects(不可恢复) /// 彻底删除:事务级联删 branches→releases→tasks→projects(不可恢复)
@@ -709,21 +715,21 @@ impl ProjectRepo {
let id = id.to_owned(); let id = id.to_owned();
tokio::task::spawn_blocking(move || { tokio::task::spawn_blocking(move || {
let mut guard = conn.blocking_lock(); 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]) 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]) 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]) 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 let affected = tx
.execute("DELETE FROM projects WHERE id = ?1", params![id]) .execute("DELETE FROM projects WHERE id = ?1", params![id])
.map_err(|e| Error::Storage(e.to_string()))?; .map_err(storage_err)?;
tx.commit().map_err(|e| Error::Storage(e.to_string()))?; tx.commit().map_err(storage_err)?;
Ok(affected > 0) Ok(affected > 0)
}) })
.await .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 guard = conn.blocking_lock();
let mut stmt = guard 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") .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 let rows = stmt
.query_map([], |row| task_from_row(row)) .query_map([], |row| task_from_row(row))
.map_err(|e| Error::Storage(e.to_string()))?; .map_err(storage_err)?;
let mut results = Vec::new(); let mut results = Vec::new();
for r in rows { for r in rows {
results.push(r.map_err(|e| Error::Storage(e.to_string()))?); results.push(r.map_err(storage_err)?);
} }
Ok(results) Ok(results)
}) })
.await .await
.map_err(|e| Error::Storage(e.to_string()))? .map_err(storage_err)?
} }
/// 软删:标记 deleted_at(进回收站,可恢复)。仅作用于未删任务,返回是否命中。 /// 软删:标记 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", "UPDATE tasks SET deleted_at = ?1, updated_at = ?1 WHERE id = ?2 AND deleted_at IS NULL",
params![now, id], params![now, id],
) )
.map_err(|e| Error::Storage(e.to_string()))?; .map_err(storage_err)?;
Ok(affected > 0) Ok(affected > 0)
}) })
.await .await
.map_err(|e| Error::Storage(e.to_string()))? .map_err(storage_err)?
} }
/// 恢复:清 deleted_at(从回收站还原)。仅作用于已删任务,返回是否命中。 /// 恢复:清 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", "UPDATE tasks SET deleted_at = NULL, updated_at = ?1 WHERE id = ?2 AND deleted_at IS NOT NULL",
params![now, id], params![now, id],
) )
.map_err(|e| Error::Storage(e.to_string()))?; .map_err(storage_err)?;
Ok(affected > 0) Ok(affected > 0)
}) })
.await .await
.map_err(|e| Error::Storage(e.to_string()))? .map_err(storage_err)?
} }
/// 原子推进任务状态(任务推进链 F-260616-02 唯一 status 写入路径) /// 原子推进任务状态(任务推进链 F-260616-02 唯一 status 写入路径)
@@ -861,22 +867,22 @@ impl TaskRepo {
}; };
let affected = guard let affected = guard
.execute(sql, params![new_status, now, id, expected]) .execute(sql, params![new_status, now, id, expected])
.map_err(|e| Error::Storage(e.to_string()))?; .map_err(storage_err)?;
if affected == 0 { if affected == 0 {
return Ok(None); return Ok(None);
} }
// 回读更新后的记录(含新 status / 累加后的 review_rounds / 新 updated_at)。 // 回读更新后的记录(含新 status / 累加后的 review_rounds / 新 updated_at)。
let mut stmt = guard 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") .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 let row = stmt
.query_row(params![id], |row| task_from_row(row)) .query_row(params![id], |row| task_from_row(row))
.optional() .optional()
.map_err(|e| Error::Storage(e.to_string()))?; .map_err(storage_err)?;
Ok(row) Ok(row)
}) })
.await .await
.map_err(|e| Error::Storage(e.to_string()))? .map_err(storage_err)?
} }
/// 列出回收站(deleted_at IS NOT NULL),按更新时间(≈删除时间)降序。对标 ProjectRepo::list_deleted。 /// 列出回收站(deleted_at IS NOT NULL),按更新时间(≈删除时间)降序。对标 ProjectRepo::list_deleted。
@@ -889,18 +895,18 @@ impl TaskRepo {
let guard = conn.blocking_lock(); let guard = conn.blocking_lock();
let mut stmt = guard 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") .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 let rows = stmt
.query_map([], |row| task_from_row(row)) .query_map([], |row| task_from_row(row))
.map_err(|e| Error::Storage(e.to_string()))?; .map_err(storage_err)?;
let mut results = Vec::new(); let mut results = Vec::new();
for r in rows { for r in rows {
results.push(r.map_err(|e| Error::Storage(e.to_string()))?); results.push(r.map_err(storage_err)?);
} }
Ok(results) Ok(results)
}) })
.await .await
.map_err(|e| Error::Storage(e.to_string()))? .map_err(storage_err)?
} }
} }
@@ -1198,15 +1204,15 @@ impl AiToolExecutionRepo {
.prepare( .prepare(
"SELECT * FROM ai_tool_executions WHERE tool_call_id = ?1 ORDER BY requested_at DESC LIMIT 1", "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 let row = stmt
.query_row(params![tid], |row| ai_tool_execution_from_row(row)) .query_row(params![tid], |row| ai_tool_execution_from_row(row))
.optional() .optional()
.map_err(|e| Error::Storage(e.to_string()))?; .map_err(storage_err)?;
Ok(row) Ok(row)
}) })
.await .await
.map_err(|e| Error::Storage(e.to_string()))? .map_err(storage_err)?
} }
/// 列出所有 status=pending 的审计行(启动重建 pending_approvals 用) /// 列出所有 status=pending 的审计行(启动重建 pending_approvals 用)
@@ -1218,18 +1224,18 @@ impl AiToolExecutionRepo {
let guard = conn.blocking_lock(); let guard = conn.blocking_lock();
let mut stmt = guard let mut stmt = guard
.prepare("SELECT * FROM ai_tool_executions WHERE status = 'pending' ORDER BY requested_at ASC") .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 let rows = stmt
.query_map([], |row| ai_tool_execution_from_row(row)) .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(); let mut results = Vec::new();
for r in rows { for r in rows {
results.push(r.map_err(|e| Error::Storage(e.to_string()))?); results.push(r.map_err(storage_err)?);
} }
Ok(results) Ok(results)
}) })
.await .await
.map_err(|e| Error::Storage(e.to_string()))? .map_err(storage_err)?
} }
} }
@@ -1281,28 +1287,28 @@ impl KnowledgeRepo {
if let Some(k) = &kind { if let Some(k) = &kind {
let mut stmt = guard 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")) .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 let rows = stmt
.query_map(params![pattern, pattern, k, limit_i], |row| knowledge_from_row(row)) .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 { for r in rows {
results.push(r.map_err(|e| Error::Storage(e.to_string()))?); results.push(r.map_err(storage_err)?);
} }
} else { } else {
let mut stmt = guard 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")) .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 let rows = stmt
.query_map(params![pattern, pattern, limit_i], |row| knowledge_from_row(row)) .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 { for r in rows {
results.push(r.map_err(|e| Error::Storage(e.to_string()))?); results.push(r.map_err(storage_err)?);
} }
} }
Ok(results) Ok(results)
}) })
.await .await
.map_err(|e| Error::Storage(e.to_string()))? .map_err(storage_err)?
} }
/// 按状态列出(审核收件箱用): 按 confidence 语义排序(high>medium>low),次按 created_at /// 按状态列出(审核收件箱用): 按 confidence 语义排序(high>medium>low),次按 created_at
@@ -1325,18 +1331,18 @@ impl KnowledgeRepo {
ELSE 0 ELSE 0
END DESC, created_at DESC", END DESC, created_at DESC",
) )
.map_err(|e| Error::Storage(e.to_string()))?; .map_err(storage_err)?;
let rows = stmt let rows = stmt
.query_map(params![status], |row| knowledge_from_row(row)) .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(); let mut results = Vec::new();
for r in rows { for r in rows {
results.push(r.map_err(|e| Error::Storage(e.to_string()))?); results.push(r.map_err(storage_err)?);
} }
Ok(results) Ok(results)
}) })
.await .await
.map_err(|e| Error::Storage(e.to_string()))? .map_err(storage_err)?
} }
/// 复用计数 +1(SQL 行级原子操作,并发安全) /// 复用计数 +1(SQL 行级原子操作,并发安全)
@@ -1351,11 +1357,11 @@ impl KnowledgeRepo {
"UPDATE knowledges SET reuse_count = reuse_count + 1, updated_at = ?1 WHERE id = ?2", "UPDATE knowledges SET reuse_count = reuse_count + 1, updated_at = ?1 WHERE id = ?2",
params![now, id], params![now, id],
) )
.map_err(|e| Error::Storage(e.to_string()))?; .map_err(storage_err)?;
Ok(affected > 0) Ok(affected > 0)
}) })
.await .await
.map_err(|e| Error::Storage(e.to_string()))? .map_err(storage_err)?
} }
/// 写入向量嵌入(BLOB = Vec<f32> 小端字节序列化) /// 写入向量嵌入(BLOB = Vec<f32> 小端字节序列化)
@@ -1372,11 +1378,11 @@ impl KnowledgeRepo {
"UPDATE knowledges SET embedding = ?1 WHERE id = ?2", "UPDATE knowledges SET embedding = ?1 WHERE id = ?2",
params![blob, id], params![blob, id],
) )
.map_err(|e| Error::Storage(e.to_string()))?; .map_err(storage_err)?;
Ok(affected > 0) Ok(affected > 0)
}) })
.await .await
.map_err(|e| Error::Storage(e.to_string()))? .map_err(storage_err)?
} }
/// 向量检索: 加载全部 published 且有 embedding 的记录,纯 Rust 余弦相似度取 top-N /// 向量检索: 加载全部 published 且有 embedding 的记录,纯 Rust 余弦相似度取 top-N
@@ -1404,18 +1410,18 @@ impl KnowledgeRepo {
.prepare(&format!( .prepare(&format!(
"SELECT {KNOWLEDGE_COLS_WITH_EMBEDDING} FROM knowledges WHERE status = 'published' AND embedding IS NOT NULL" "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 let rows = stmt
.query_map([], |row| { .query_map([], |row| {
let rec = knowledge_from_row(row)?; let rec = knowledge_from_row(row)?;
let blob: Vec<u8> = row.get("embedding")?; let blob: Vec<u8> = row.get("embedding")?;
Ok((rec, blob)) Ok((rec, blob))
}) })
.map_err(|e| Error::Storage(e.to_string()))?; .map_err(storage_err)?;
let mut scored: Vec<(KnowledgeRecord, f32)> = Vec::new(); let mut scored: Vec<(KnowledgeRecord, f32)> = Vec::new();
for r in rows { 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); let emb = blob_to_f32s(&blob);
// 维度不匹配(换过 embedding 模型的旧向量)跳过 // 维度不匹配(换过 embedding 模型的旧向量)跳过
if emb.len() != query_vec.len() { if emb.len() != query_vec.len() {
@@ -1429,7 +1435,7 @@ impl KnowledgeRepo {
Ok(scored) Ok(scored)
}) })
.await .await
.map_err(|e| Error::Storage(e.to_string()))? .map_err(storage_err)?
} }
/// 列出非归档知识(全部 status != 'archived'),按 confidence 语义排序 /// 列出非归档知识(全部 status != 'archived'),按 confidence 语义排序
@@ -1449,18 +1455,18 @@ impl KnowledgeRepo {
ELSE 0 ELSE 0
END DESC, created_at DESC", END DESC, created_at DESC",
) )
.map_err(|e| Error::Storage(e.to_string()))?; .map_err(storage_err)?;
let rows = stmt let rows = stmt
.query_map([], |row| knowledge_from_row(row)) .query_map([], |row| knowledge_from_row(row))
.map_err(|e| Error::Storage(e.to_string()))?; .map_err(storage_err)?;
let mut results = Vec::new(); let mut results = Vec::new();
for r in rows { for r in rows {
results.push(r.map_err(|e| Error::Storage(e.to_string()))?); results.push(r.map_err(storage_err)?);
} }
Ok(results) Ok(results)
}) })
.await .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 guard = conn.blocking_lock();
let mut stmt = guard 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") .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 let rows = stmt
.query_map(params![limit_i], |row| knowledge_from_row(row)) .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(); let mut results = Vec::new();
for r in rows { for r in rows {
results.push(r.map_err(|e| Error::Storage(e.to_string()))?); results.push(r.map_err(storage_err)?);
} }
Ok(results) Ok(results)
}) })
.await .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 guard = conn.blocking_lock();
let mut stmt = guard 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") .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 let rows = stmt
.query_map(params![knowledge_id], |row| knowledge_event_from_row(row)) .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(); let mut results = Vec::new();
for r in rows { for r in rows {
results.push(r.map_err(|e| Error::Storage(e.to_string()))?); results.push(r.map_err(storage_err)?);
} }
Ok(results) Ok(results)
}) })
.await .await
.map_err(|e| Error::Storage(e.to_string()))? .map_err(storage_err)?
} }
/// 按知识 ID + event_type 查询最近 N 条(如引用记录翻页: type=referenced) /// 按知识 ID + event_type 查询最近 N 条(如引用记录翻页: type=referenced)
@@ -1545,18 +1551,18 @@ impl KnowledgeEventsRepo {
let guard = conn.blocking_lock(); let guard = conn.blocking_lock();
let mut stmt = guard 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") .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 let rows = stmt
.query_map(params![knowledge_id, event_type, limit_i], |row| knowledge_event_from_row(row)) .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(); let mut results = Vec::new();
for r in rows { for r in rows {
results.push(r.map_err(|e| Error::Storage(e.to_string()))?); results.push(r.map_err(storage_err)?);
} }
Ok(results) Ok(results)
}) })
.await .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", "UPDATE ai_conversations SET messages = '[]', prompt_tokens = 0, completion_tokens = 0, updated_at = ?1 WHERE id = ?2",
params![now, id], params![now, id],
) )
.map_err(|e| Error::Storage(e.to_string()))?; .map_err(storage_err)?;
Ok(affected > 0) Ok(affected > 0)
}) })
.await .await
.map_err(|e| Error::Storage(e.to_string()))? .map_err(storage_err)?
} }
/// 设置归档标记(仅改 archived,不动 updated_at) /// 设置归档标记(仅改 archived,不动 updated_at)
@@ -1599,11 +1605,11 @@ impl AiConversationRepo {
"UPDATE ai_conversations SET archived = ?1 WHERE id = ?2", "UPDATE ai_conversations SET archived = ?1 WHERE id = ?2",
params![if archived { 1 } else { 0 }, id], params![if archived { 1 } else { 0 }, id],
) )
.map_err(|e| Error::Storage(e.to_string()))?; .map_err(storage_err)?;
Ok(affected > 0) Ok(affected > 0)
}) })
.await .await
.map_err(|e| Error::Storage(e.to_string()))? .map_err(storage_err)?
} }
} }

View File

@@ -32,64 +32,30 @@ pub fn run(conn: &Connection) -> Result<()> {
) )
.unwrap_or(0); .unwrap_or(0);
if current_version < 1 { // 迁移步骤链: 顺序执行,跳过已应用的版本(current_version < N 才跑)。
migrate_v1(conn)?; // 新增版本时,在此数组追加一项 (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 { for (version, migrate_fn) in steps {
migrate_v2(conn)?; if current_version < version {
migrate_fn(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)?;
} }
Ok(()) Ok(())