优化: 边界加固(AI loop竞态根治+数据/审批/反馈/并发/错误分类)

AI loop 竞态(P0):per-conv epoch/owner token + 存活心跳治 force_send 双loop + stop 3s兜底误判;旧loop stale 全跳过(guard/emit/save)

agentic 收尾(A2-B8):Fatal 退出落库user消息(镜像Exhausted)+ 入口早退补save + usage is_estimated 打标 + emit_ai_completed_once 单点收敛清审批残留

聊天清理(A2-B9):clearChat 先停loop→DB单事务→内存清(clear_conversation_atomic)+ 前端错误气泡

循环并发(A2-B11):三态 ProviderAcquire(NotConfigured/Acquired/Exhausted)+ 候选循环非阻塞+防抖3次饱和降级+单测

错误分类(A2-B12):stream error帧接入 classify_status_or_class + 关键词保守降级 + 7单测

数据(G1.2/G1.4):purge_with_descendants 级联补全(11表单事务+存在性守卫)+ move_task_queue 单事务收口(两调用方共用)

git只读(G3.1):run_git_status/diff/log success判定(exit_code差异语义,失败结构化{success:false,error})

安全(G5.2/G5.6):create_project 目录Err+name校验 + module.rs 路径遍历DRY(分段匹配修a..b.rs误伤)

幂等(V2/V32):裸ALTER全守卫化 + v1..v40全链重跑幂等测试(16过)

附:remote_bridge await 临时引用修(E0716)+ agentic emit 收敛 E0716 app_state 绑定修
This commit is contained in:
lxy
2026-08-05 22:10:32 +08:00
parent ec9f0bf1ea
commit 5667da6cf4
16 changed files with 1696 additions and 215 deletions
@@ -524,6 +524,43 @@ impl AiConversationRepo {
.map_err(storage_err)?
}
/// 清空对话消息内容(单事务原子:ai_conversations.messages 置 '[]' + ai_messages 表全删)。
///
/// A2-B9(G3.2 clearChat 裁决):原 `clear_messages` + `delete_range` 两条独立 DB 写非原子,
/// DB 失败会致 messages JSON 列与 ai_messages 表不一致(如仅一条成功)。本方法一次 transaction
/// 覆盖两条写(① UPDATE ai_conversations 置空消息 + 清零 token;② DELETE ai_messages 该 conv
/// 全部行),成功全成功 / 失败回滚全失败。供 `ai_chat_clear` 先停 loop 再单事务清空。
///
/// 对话壳保留(侧栏仍可见,可继续在该对话内聊);返回 Ok(())——调用方只关心成功与否
/// (对齐 replace_conversation 语义,不返回受影响行数)。
pub async fn clear_conversation_atomic(&self, id: &str) -> Result<()> {
let conn = self.conn.clone();
let id = id.to_owned();
let now = now_millis_str();
tokio::task::spawn_blocking(move || -> Result<()> {
let mut guard = conn.blocking_lock();
let tx = guard.transaction().map_err(storage_err)?;
{
// ① ai_conversations.messages 置空 + token 清零(对话壳保留)
tx.execute(
"UPDATE ai_conversations SET messages = '[]', prompt_tokens = 0, completion_tokens = 0, updated_at = ?1 WHERE id = ?2",
params![now, id],
)
.map_err(storage_err)?;
// ② ai_messages 表全删(等价 delete_range min_seq=0 max=None:seq 恒 >= 0)
tx.execute(
"DELETE FROM ai_messages WHERE conversation_id = ?1",
params![id],
)
.map_err(storage_err)?;
}
tx.commit().map_err(storage_err)?;
Ok(())
})
.await
.map_err(storage_err)?
}
/// 设置归档标记(仅改 archived,不动 updated_at)
///
/// 区别于 update_field(后者强制 SET updated_at=now,会把归档/取消归档误判为内容更新,
@@ -586,6 +623,37 @@ impl AiConversationRepo {
.await
.map_err(storage_err)?
}
/// G1.3: 删除对话 + 其全部 ai_messages 子行(单事务原子)。
///
/// 背景:原宏生成 `delete` 只删 ai_conversations 主行,而 ai_messages 表无外键级联
/// (conversation_id 仅普通索引),子行孤儿累积。本方法在同一事务内**先删子行
/// (ai_messages)再删主行(ai_conversations)**,要么全删要么全不删。
///
/// 顺序注意:先删数据再摘内存(命令层 per_conv.remove 在其后),防后台在途
/// save_conversation 在删主行后把孤儿消息写回复活。与 save_conversation 共享同一
/// conn(Mutex),事务原子性保证删除期间无中间态(半删半留)。
pub async fn delete_with_messages(&self, id: &str) -> Result<bool> {
let conn = self.conn.clone();
let id = id.to_owned();
tokio::task::spawn_blocking(move || {
let mut guard = conn.blocking_lock();
let tx = guard.transaction().map_err(storage_err)?;
// 先删子行(ai_messages)再删主行(ai_conversations),单事务原子
tx.execute(
"DELETE FROM ai_messages WHERE conversation_id = ?1",
params![id],
)
.map_err(storage_err)?;
let conv_affected = tx
.execute("DELETE FROM ai_conversations WHERE id = ?1", params![id])
.map_err(storage_err)?;
tx.commit().map_err(storage_err)?;
Ok(conv_affected > 0)
})
.await
.map_err(storage_err)?
}
}
// ============================================================
+70 -7
View File
@@ -414,22 +414,85 @@ impl ProjectRepo {
.map_err(storage_err)?
}
/// 彻底删除:事务级联删 branches→releases→tasks→projects(不可恢复)
/// 彻底删除:事务级联删全部关联子表→projects(不可恢复)
///
/// SQLite 已开 PRAGMA foreign_keys=ON 但表无 ON DELETE CASCADE,ALTER 改不了 FK 约束,
/// 故应用层级联:子表先于父表删,单事务保证一致性。
///
/// **级联范围**(G1.2 补全):V1-V35 全部 `REFERENCES projects(id)` 表——
/// branches / releases / tasks / workflow_executions(V2) + 知识图谱 V29-V35 新增
/// task_links / project_events / project_services / project_modules / module_dependencies,
/// 以及经它们间接到项目的 node_executions(REFERENCES workflow_executions)/
/// module_dependencies(REFERENCES project_modules)。此前只删四表,带 modules/事件/服务的
/// 工程 purge 在 foreign_keys=ON 下必 FK 违例回滚(latent bug)。
///
/// **删除顺序 = FK 拓扑 最深子表→父**(先删被引用方,防 foreign_keys=ON 下 FK 违例):
/// node_executions → task_links → branches → module_dependencies → workflow_executions
/// → tasks → project_events → project_services → project_modules → releases → projects。
///
/// 表存在性守卫:逐表先查 `sqlite_master` 存在才 DELETE(防老库缺 V29-V35 新表时
/// `no such table` 中断整个事务回滚)。
pub async fn purge_with_descendants(&self, id: &str) -> Result<bool> {
let conn = self.conn.clone();
let id = id.to_owned();
tokio::task::spawn_blocking(move || {
let mut guard = conn.blocking_lock();
let tx = guard.transaction().map_err(storage_err)?;
tx.execute("DELETE FROM branches WHERE project_id = ?1", params![id])
.map_err(storage_err)?;
tx.execute("DELETE FROM releases WHERE project_id = ?1", params![id])
.map_err(storage_err)?;
tx.execute("DELETE FROM tasks WHERE project_id = ?1", params![id])
.map_err(storage_err)?;
// (表名, 删除 SQL)。表名/SQL 均为编译期常量,无注入风险。
// node_executions / task_links 无 project_id 列,经其父表子查询收敛到本项目;
// workflow_executions 兼删 task_id 命中(project_id 可空,补齐 task_id 关联路径)。
const CASCADE: &[(&str, &str)] = &[
// node_executions REFERENCES workflow_executions(id):先于 workflow_executions 删
(
"node_executions",
"DELETE FROM node_executions WHERE workflow_id IN \
(SELECT id FROM workflow_executions WHERE project_id = ?1 \
OR task_id IN (SELECT id FROM tasks WHERE project_id = ?1))",
),
// task_links REFERENCES tasks(id):先于 tasks 删(source/target 任一端属本项目)
(
"task_links",
"DELETE FROM task_links WHERE source_id IN \
(SELECT id FROM tasks WHERE project_id = ?1) OR target_id IN \
(SELECT id FROM tasks WHERE project_id = ?1)",
),
// branches REFERENCES tasks(id) + projects(id):先于 tasks 删
("branches", "DELETE FROM branches WHERE project_id = ?1"),
// module_dependencies REFERENCES project_modules(id) + projects(id):先于 project_modules 删
("module_dependencies", "DELETE FROM module_dependencies WHERE project_id = ?1"),
// workflow_executions REFERENCES projects(id)(可空,兼删 task_id 关联)
(
"workflow_executions",
"DELETE FROM workflow_executions WHERE project_id = ?1 \
OR task_id IN (SELECT id FROM tasks WHERE project_id = ?1)",
),
// tasks REFERENCES projects(id)(自引用 parent_id 同语句删,单语句末校验不违例)
("tasks", "DELETE FROM tasks WHERE project_id = ?1"),
// 以下直接 REFERENCES projects(id),无子表依赖,顺序任意
("project_events", "DELETE FROM project_events WHERE project_id = ?1"),
("project_services", "DELETE FROM project_services WHERE project_id = ?1"),
("project_modules", "DELETE FROM project_modules WHERE project_id = ?1"),
("releases", "DELETE FROM releases WHERE project_id = ?1"),
];
for (table, sql) in CASCADE {
// 表存在性守卫:老库可能缺 V29-V35 新表,缺表时 DELETE 报 no such table
// 中断事务回滚;探测存在才删,防老库 purge 失败。
let exists: bool = tx
.query_row(
"SELECT 1 FROM sqlite_master WHERE type = 'table' AND name = ?1",
params![*table],
|_| Ok(true),
)
.optional()
.map_err(storage_err)?
.unwrap_or(false);
if exists {
tx.execute(sql, params![id]).map_err(storage_err)?;
}
}
let affected = tx
.execute("DELETE FROM projects WHERE id = ?1", params![id])
.map_err(storage_err)?;
+220 -1
View File
@@ -567,7 +567,8 @@ impl TaskRepo {
///
/// 防护:
/// - 方法名 `set_status_for_aggregation` 显式表明语义,非通用 setter,防误用。
/// - 调用方(commands::task::recompute_parent_status)负责聚合规则计算,本方法只落库。
/// - 调用方(df-nodes task_advance_node::recompute_parent_status,2026-08-04 下沉共享层)负责
/// 聚合规则计算,本方法只落库。
/// - 不动 review_rounds(父任务不执行工作流,无 review 退回语义)。
///
/// 返回是否命中(父任务不存在/已删 → false)。
@@ -591,6 +592,120 @@ impl TaskRepo {
.map_err(storage_err)?
}
/// 跨池移动任务(单事务原子:读当前 → 一致性联动 status → 写 queue + status)。
///
/// 整合原 IPC 层(src-tauri commands::task::move_task_queue)与 AI 工具
/// (src-tauri commands::ai::tools::task_graph)各自的两段式
/// (`update_field` queue + `set_status_for_aggregation` status)为**单一 repo 方法**:
/// 一次 `blocking_lock` 内 get + update queue + update status 于**同一 transaction**,
/// 杜绝「queue 已改、status 未改」的非原子中间态(G1.4)。两调用方共用本方法防漂移。
///
/// **不能 naive 在 tx 内复用既有 repo 方法**(各自内部 `blocking_lock`,std Mutex
/// 非重入必死锁),故本方法用裸 SQL 在事务内完成全部读写。
///
/// 一致性约束联动(设计 §2.1):
/// - queue=done → status 强制=done(池完成即任务完成)
/// - queue=backlog → status 强制=todo(需求池任务尚未开始)
/// - queue=active → status 若不在 {in_progress,in_review,testing} 则强制=in_progress
/// - queue=todo → status 非 todo 则强制=todo(待办池任务尚未开始)
/// - queue=decision → status 不变(待决策池保留执行态,暂停推进不重置)
///
/// 父任务(容器模型)也可 move_task_queue(其 status 由聚合规则 recompute_parent_status
/// 在子任务推进时重算,本方法仅满足一致性约束联动,不动 review_rounds)。
///
/// - queue 白名单校验(bad_queue 防进 `_ => unreachable!` match)收口在方法内;
/// - `deleted_at IS NULL` 收口:软删回收站任务不可 move(与 set_status_for_aggregation
/// 语义一致,返回 None);
/// - 返回:更新后的 TaskRecord(Some);任务不存在/已软删 → None(调用方据此报「任务不存在」)。
pub async fn move_task_queue(
&self,
id: &str,
new_queue: &str,
) -> Result<Option<TaskRecord>> {
// queue 白名单校验(对标 commands::task::validate_queue,防非法值进 match unreachable)。
// 常量与联动规则集中在 Repo 层,commands/task.rs 与 ai/tools/task_graph.rs 两调用方
// 共用同一方法(防漂移),不再各自实现。
const TASK_QUEUE_VALUES: &[&str] = &["backlog", "todo", "decision", "active", "done"];
const ACTIVE_OK_STATUSES: &[&str] = &["in_progress", "in_review", "testing"];
if !TASK_QUEUE_VALUES.contains(&new_queue) {
return Err(df_types::error::Error::Validation(format!(
"非法 queue 值 {:?},合法值: {:?}",
new_queue, TASK_QUEUE_VALUES
)));
}
let conn = self.conn.clone();
let id = id.to_owned();
let new_queue = new_queue.to_owned();
let now = now_millis_str();
// 显式列出全部 18 列(同 from_row 消费列,不 SELECT deleted_at:
// TaskRecord 不带该字段,取了 from_row 会因未知列报错)。
const TASK_COLS: &str = "id, project_id, title, description, status, priority, branch_name, \
assignee, workflow_def_id, base_branch, review_rounds, output_json, \
idea_id, queue, parent_id, content_json, created_at, updated_at";
tokio::task::spawn_blocking(move || {
let mut guard = conn.blocking_lock();
let tx = guard.transaction().map_err(storage_err)?;
// 1. 读当前(取 status 做一致性联动决策;deleted_at IS NULL 收口软删任务不可 move)。
let current: Option<TaskRecord> = {
let mut stmt = tx
.prepare(&format!(
"SELECT {TASK_COLS} FROM tasks WHERE id = ?1 AND deleted_at IS NULL"
))
.map_err(storage_err)?;
stmt.query_row(params![id], |row| task_from_row(row))
.optional()
.map_err(storage_err)?
};
let Some(current) = current else {
return Ok(None);
};
// 2. 一致性联动:根据 new_queue 决定 status 是否需调整(设计 §2.1)。
let new_status = match new_queue.as_str() {
"done" => "done".to_string(),
"backlog" => "todo".to_string(),
"active" => {
if ACTIVE_OK_STATUSES.contains(&current.status.as_str()) {
current.status.as_str().to_string() // 已在执行中三态,保留
} else {
"in_progress".to_string() // 否则强制进 in_progress
}
}
"todo" => "todo".to_string(),
"decision" => current.status.as_str().to_string(), // 保留执行态
_ => unreachable!("queue 白名单已收口"),
};
// 3. 同一事务内写 queue + status(与 current 不同才写,避免无谓 updated_at 抖动)。
let queue_changed = current.queue != new_queue;
let status_changed = current.status.as_str() != new_status;
if queue_changed || status_changed {
tx.execute(
"UPDATE tasks SET queue = ?1, status = ?2, updated_at = ?3 WHERE id = ?4",
params![new_queue, new_status, now, id],
)
.map_err(storage_err)?;
}
// 4. 回读最新记录返回。
let updated: Option<TaskRecord> = {
let mut stmt = tx
.prepare(&format!("SELECT {TASK_COLS} FROM tasks WHERE id = ?1"))
.map_err(storage_err)?;
stmt.query_row(params![id], |row| task_from_row(row))
.optional()
.map_err(storage_err)?
};
tx.commit().map_err(storage_err)?;
Ok(updated)
})
.await
.map_err(storage_err)?
}
/// 列出回收站(deleted_at IS NOT NULL),按更新时间(≈删除时间)降序。对标 ProjectRepo::list_deleted。
///
/// 注:按项目列活跃任务走 list_active_by_project(SQL 下推 project_id),
@@ -866,4 +981,108 @@ mod tests {
let ok = repo.set_status_for_aggregation("ghost", "done").await.unwrap();
assert!(!ok, "不存在的任务应返回 false");
}
// ============================================================
// move_task_queue 单事务跨池移动(G1.4:读当前 → 联动 status → 写 queue+status 原子)
// ============================================================
/// 移动后断言 queue/status 双落地(单事务原子,非两段式独立写)。
#[tokio::test]
async fn move_task_queue_done_sets_queue_and_status_atomically() {
let repo = setup().await;
repo.insert(trec_full("t1", "active", None, TaskStatus::InProgress))
.await
.unwrap();
let updated = repo.move_task_queue("t1", "done").await.unwrap().unwrap();
assert_eq!(updated.queue, "done", "queue 应改为 done");
assert_eq!(updated.status.as_str(), "done", "status 应联动强制 done");
}
#[tokio::test]
async fn move_task_queue_backlog_forces_todo_status() {
let repo = setup().await;
repo.insert(trec_full("t1", "active", None, TaskStatus::InProgress))
.await
.unwrap();
let updated = repo.move_task_queue("t1", "backlog").await.unwrap().unwrap();
assert_eq!(updated.queue, "backlog");
assert_eq!(updated.status.as_str(), "todo", "backlog 池 status 强制 todo");
}
#[tokio::test]
async fn move_task_queue_active_preserves_executing_status() {
// active 池:已在执行中三态(in_progress)则保留,不重置
let repo = setup().await;
repo.insert(trec_full("t1", "backlog", None, TaskStatus::InProgress))
.await
.unwrap();
let updated = repo.move_task_queue("t1", "active").await.unwrap().unwrap();
assert_eq!(updated.queue, "active");
assert_eq!(updated.status.as_str(), "in_progress", "active 池保留执行中三态");
}
#[tokio::test]
async fn move_task_queue_active_forces_in_progress_when_idle() {
// active 池:非执行中三态(todo)→ 强制 in_progress
let repo = setup().await;
repo.insert(trec_full("t1", "todo", None, TaskStatus::Todo))
.await
.unwrap();
let updated = repo.move_task_queue("t1", "active").await.unwrap().unwrap();
assert_eq!(updated.queue, "active");
assert_eq!(updated.status.as_str(), "in_progress", "非执行态进 active 强制 in_progress");
}
#[tokio::test]
async fn move_task_queue_decision_preserves_status() {
// decision 池保留当前执行态(不重置)
let repo = setup().await;
repo.insert(trec_full("t1", "active", None, TaskStatus::InReview))
.await
.unwrap();
let updated = repo.move_task_queue("t1", "decision").await.unwrap().unwrap();
assert_eq!(updated.queue, "decision");
assert_eq!(updated.status.as_str(), "in_review", "decision 池保留执行态");
}
#[tokio::test]
async fn move_task_queue_soft_deleted_returns_none() {
// 软删回收站任务不可 move(deleted_at IS NULL 收口,与 set_status_for_aggregation 一致)
let repo = setup().await;
repo.insert(trec_full("t1", "todo", None, TaskStatus::Todo))
.await
.unwrap();
repo.soft_delete("t1").await.unwrap();
let res = repo.move_task_queue("t1", "done").await.unwrap();
assert!(res.is_none(), "软删任务 move 应返回 None");
}
#[tokio::test]
async fn move_task_queue_nonexistent_returns_none() {
let repo = setup().await;
let res = repo.move_task_queue("ghost", "done").await.unwrap();
assert!(res.is_none(), "不存在的任务 move 应返回 None");
}
#[tokio::test]
async fn move_task_queue_invalid_queue_rejected() {
let repo = setup().await;
repo.insert(trec_full("t1", "todo", None, TaskStatus::Todo))
.await
.unwrap();
let err = repo.move_task_queue("t1", "bogus").await.unwrap_err();
assert!(matches!(err, df_types::error::Error::Validation(_)), "非法 queue 应拒绝");
}
#[tokio::test]
async fn move_task_queue_same_queue_noop() {
// 同池 no-op:queue/status 均不变,updated_at 不抖动(updated 记录仍正常返回)
let repo = setup().await;
repo.insert(trec_full("t1", "todo", None, TaskStatus::Todo))
.await
.unwrap();
let updated = repo.move_task_queue("t1", "todo").await.unwrap().unwrap();
assert_eq!(updated.queue, "todo");
assert_eq!(updated.status.as_str(), "todo");
}
}