Files
DevFlow/crates/df-storage/src/crud/task_repo.rs
绝尘 c4b02b5370 修复: 消息溯源P2切读确认完成 + 走查报告高优项
- P2 切读方案文档更新(确认批次B已上线)
- query LIKE 通配符转义(project/idea/task/knowledge 4 repo)
- dag.rs deep_merge null 覆盖全局配置
- dag.rs 通配 match 展开显式变体
- INDEX.md 补漏 10 个文档索引
- ARCHITECTURE.md 删除与新文档逐字重复
- coordinator.rs 加 #[deprecated] 编译守卫
- AiChat.vue 空值传播加 console.warn
2026-06-28 05:06:51 +08:00

816 lines
38 KiB
Rust
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
//! 任务域 Repo:TaskRepo(含 advance_status_atomic 状态机收口)
use std::sync::Arc;
use rusqlite::{params, Connection, OptionalExtension, Row};
use serde::{Deserialize, Serialize};
use tokio::sync::Mutex;
use df_types::error::Result;
use df_types::types::TaskStatus;
use crate::db::Database;
use crate::models::TaskRecord;
use super::impl_repo;
use super::{now_millis_str, storage_err, validate_column_name};
// ============================================================
// from_row 辅助函数
// ============================================================
fn task_from_row(row: &Row<'_>) -> std::result::Result<TaskRecord, rusqlite::Error> {
Ok(TaskRecord {
id: row.get("id")?,
project_id: row.get("project_id")?,
title: row.get("title")?,
description: row.get("description")?,
status: {
let s: String = row.get("status")?;
TaskStatus::from_db_str(&s).unwrap_or_default()
},
priority: row.get("priority")?,
branch_name: row.get("branch_name")?,
assignee: row.get("assignee")?,
workflow_def_id: row.get("workflow_def_id")?,
base_branch: row.get("base_branch")?,
review_rounds: row.get("review_rounds")?,
output_json: row.get("output_json")?,
idea_id: row.get("idea_id")?,
// 知识图谱 Phase 1 V29 三列(queue/parent_id/content_json),15→18 列同步之一。
queue: row.get("queue")?,
parent_id: row.get("parent_id")?,
content_json: row.get("content_json")?,
created_at: row.get("created_at")?,
updated_at: row.get("updated_at")?,
})
}
// ============================================================
// 任务列表查询入参(F-260621-02 查询维度补全)
// ============================================================
/// 任务列表动态查询入参。全可选,空 query = 等价当前全量行为(向后兼容)。
///
/// 设计对齐方案文档(docs/02-架构设计/专项设计/查询能力补全方案-2026-06-21.md §4.1):
/// ① 字段全 Option,旧调用方不传 / 传空 → 等价 list_active(全量未删),零破坏;
/// ② 复用 KnowledgeRepo::search 的动态 WHERE 拼接模式(if-let 分支拼 SQL + 参数绑定);
/// ③ order_by 走白名单(TASK_ORDER_BY_WHITELIST)防 SQL 注入,对齐
/// impl_repo! 宏 validate_column_name 的白名单防注入思路;
/// ④ limit/offset 钳制上限(limit.min(500))防滥用,对齐 conversation_repo::list_recent
/// 的 limit.min(200)(任务场景放宽至 500,数据量 ~15 任务)。
///
/// 当前仅 status(P1 下沉)+ keyword(P2 LIKE)在视图链路使用;project_id/priority/assignee/
/// order_by/limit/offset 为基建就绪(视图暂不用,P3 排序分页待数据量增长)。
///
/// 知识图谱 Phase 1 V29(对标设计 §2.1)新增 queue / parent_id 两个筛选维度:
/// - `queue`:管理池筛选(backlog/todo/decision/active/done),看板视图按池分列的数据源;
/// - `parent_id`:父任务筛选(Some(id)=某父的子任务;特殊语义 None 仅叶子 vs 全部 由调用方拼)。
#[derive(Debug, Clone, Default, Serialize, Deserialize)]
pub struct TaskQuery {
/// 项目 ID 过滤(SQL 下推,命中 idx_tasks_project_id)
pub project_id: Option<String>,
/// 状态过滤(P1 下沉,命中 idx_tasks_status)。值集由上层 TaskStatus::is_valid 兜底校验。
pub status: Option<String>,
/// 优先级过滤(0=critical..3=low)。P3 基建就绪,视图暂不用。
pub priority: Option<i32>,
/// 负责人过滤。P3 基建就绪,视图暂不用。
pub assignee: Option<String>,
/// 关键词搜索(P2):title/description LIKE %kw%,对齐知识库 search 的 LIKE 模式。
pub keyword: Option<String>,
/// 管理池过滤(知识图谱 Phase 1 V29):backlog/todo/decision/active/done。
/// 看板视图按池分列的数据源。值合法性由上层 queue 语义校验兜底(非法值 DB 无匹配返回空)。
#[serde(default)]
pub queue: Option<String>,
/// 父任务 ID 过滤(知识图谱 Phase 1 V29):Some(id) = 查某父任务的子任务;
/// 查叶子任务(parent_id IS NULL)由专用方法 get_children 之外的语义决定,本字段只做等值匹配。
#[serde(default)]
pub parent_id: Option<String>,
/// 排序字段(白名单 created_at/updated_at/priority/status,降序)。P3 基建就绪。
pub order_by: Option<String>,
/// 分页上限(钳制 ≤500)。P3 基建就绪。
pub limit: Option<u32>,
/// 分页偏移。P3 基建就绪。
pub offset: Option<u32>,
}
/// order_by 白名单:只允许这些列名拼进 SQL(防注入,列名不可参数化只能白名单)。
/// 对齐 impl_repo! 宏 validate_column_name 的「白名单先于拼接校验」防注入思路。
/// 注:created_at 是默认值(白名单含它,query 未指定时用 created_at DESC)。
const TASK_ORDER_BY_WHITELIST: &[&str] = &["created_at", "updated_at", "priority", "status"];
/// 校验 order_by 列名在白名单内,否则返回 Err(防 SQL 注入)。
fn validate_order_by(field: &str) -> Result<()> {
if TASK_ORDER_BY_WHITELIST.contains(&field) {
Ok(())
} else {
Err(df_types::error::Error::Storage(format!(
"非法 order_by 字段名: {},合法值: {:?}",
field, TASK_ORDER_BY_WHITELIST
)))
}
}
// ============================================================
// Repo 实现
// ============================================================
impl_repo!(
/// 任务表 CRUD
TaskRepo,
TaskRecord,
"tasks",
from_row => |row| task_from_row(row),
insert => |conn, rec| {
conn.execute(
"INSERT INTO tasks (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)
VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, ?10, ?11, ?12, ?13, ?14, ?15, ?16, ?17, ?18)",
params![
rec.id, rec.project_id, rec.title, rec.description, rec.status.as_str(), rec.priority,
rec.branch_name, rec.assignee, rec.workflow_def_id, rec.base_branch,
rec.review_rounds, rec.output_json, rec.idea_id,
rec.queue, rec.parent_id, rec.content_json,
rec.created_at, rec.updated_at
],
)
},
update => |conn, rec| {
conn.execute(
"UPDATE tasks SET project_id = ?1, title = ?2, description = ?3, status = ?4, priority = ?5, branch_name = ?6, assignee = ?7, workflow_def_id = ?8, base_branch = ?9, review_rounds = ?10, output_json = ?11, idea_id = ?12, queue = ?13, parent_id = ?14, content_json = ?15, updated_at = ?16 WHERE id = ?17",
params![
rec.project_id, rec.title, rec.description, rec.status.as_str(), rec.priority,
rec.branch_name, rec.assignee, rec.workflow_def_id, rec.base_branch,
rec.review_rounds, rec.output_json, rec.idea_id,
rec.queue, rec.parent_id, rec.content_json,
rec.updated_at, rec.id
],
)
}
);
impl TaskRepo {
/// 列出未删除任务(deleted_at IS NULL)— 对标 ProjectRepo::list_active
///
/// 显式列出全部 18 个 TaskRecord 列名(同 ProjectRepo::list_active 写法),
/// 不 SELECT deleted_at:TaskRecord 不带该字段,取了 from_row 会因未知列报错。
pub async fn list_active(&self) -> Result<Vec<TaskRecord>> {
let conn = self.conn.clone();
tokio::task::spawn_blocking(move || {
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, output_json, idea_id, queue, parent_id, content_json, created_at, updated_at FROM tasks WHERE deleted_at IS NULL ORDER BY created_at DESC")
.map_err(storage_err)?;
let rows = stmt
.query_map([], |row| task_from_row(row))
.map_err(storage_err)?;
let mut results = Vec::new();
for r in rows {
results.push(r.map_err(storage_err)?);
}
Ok(results)
})
.await
.map_err(storage_err)?
}
/// 软删:标记 deleted_at(进回收站,可恢复)。仅作用于未删任务,返回是否命中。
/// 对标 ProjectRepo::soft_delete。
pub async fn soft_delete(&self, id: &str) -> Result<bool> {
let conn = self.conn.clone();
let id = id.to_owned();
let now = now_millis_str();
tokio::task::spawn_blocking(move || {
let guard = conn.blocking_lock();
let affected = guard
.execute(
"UPDATE tasks SET deleted_at = ?1, updated_at = ?1 WHERE id = ?2 AND deleted_at IS NULL",
params![now, id],
)
.map_err(storage_err)?;
Ok(affected > 0)
})
.await
.map_err(storage_err)?
}
/// 恢复:清 deleted_at(从回收站还原)。仅作用于已删任务,返回是否命中。
/// 对标 ProjectRepo::restore。
pub async fn restore(&self, id: &str) -> Result<bool> {
let conn = self.conn.clone();
let id = id.to_owned();
let now = now_millis_str();
tokio::task::spawn_blocking(move || {
let guard = conn.blocking_lock();
let affected = guard
.execute(
"UPDATE tasks SET deleted_at = NULL, updated_at = ?1 WHERE id = ?2 AND deleted_at IS NOT NULL",
params![now, id],
)
.map_err(storage_err)?;
Ok(affected > 0)
})
.await
.map_err(storage_err)?
}
/// 原子推进任务状态(任务推进链 F-260616-02 唯一 status 写入路径)
///
/// 下沉 SQL `WHERE id=? AND status=?expected` 做 CAS(Compare-And-Swap)防 TOCTOU:
/// 并发推进/旁路修改若已改 status,affected_rows==0,本方法返回 None,调用方
/// (task_advance_node)据此报「状态已变,推进中止」。`review_rounds` 不进通用
/// update_field 白名单(收口:仅本方法可改 status 与 review_rounds)。
///
/// - `expected`:调用方读取的当前 status(状态机校验时的 from),CAS 前置。
/// - `new_status`:目标 status(状态机 can_transition 已校验合法)。
/// - `bump_rounds`:退回转换(in_review→in_progress / testing→in_review)传 true,
/// 一并 `review_rounds = review_rounds + 1`(同 UPDATE 原子,避免读改写竞争)。
/// 前向推进 / 进出 blocked / 进 cancelled 传 false,不动 review_rounds。
///
/// 返回:成功推进返回更新后的 TaskRecord;affected==0(状态已变/任务不存在)返回 None。
pub async fn advance_status_atomic(
&self,
id: &str,
expected: &str,
new_status: &str,
bump_rounds: bool,
) -> Result<Option<TaskRecord>> {
let conn = self.conn.clone();
let id = id.to_owned();
let expected = expected.to_owned();
let new_status = new_status.to_owned();
let now = now_millis_str();
tokio::task::spawn_blocking(move || {
let guard = conn.blocking_lock();
// CAS:WHERE id AND status=expected 锁定当前态;affected==0 即并发已改动。
// deleted_at IS NULL:回收站任务(soft_delete 设了 deleted_at)CAS 必败→affected=0,
// 返回 None,杜绝回收站任务被推进(D-02 软删语义收口,一处关闭)。
let sql = if bump_rounds {
"UPDATE tasks SET status = ?1, review_rounds = review_rounds + 1, updated_at = ?2 \
WHERE id = ?3 AND status = ?4 AND deleted_at IS NULL"
} else {
"UPDATE tasks SET status = ?1, updated_at = ?2 \
WHERE id = ?3 AND status = ?4 AND deleted_at IS NULL"
};
let affected = guard
.execute(sql, params![new_status, now, id, expected])
.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, output_json, idea_id, queue, parent_id, content_json, created_at, updated_at FROM tasks WHERE id = ?1")
.map_err(storage_err)?;
let row = stmt
.query_row(params![id], |row| task_from_row(row))
.optional()
.map_err(storage_err)?;
Ok(row)
})
.await
.map_err(storage_err)?
}
/// 按项目列出未删除任务(deleted_at IS NULL AND project_id = ?),按创建时间降序。
///
/// SQL 下推 project_id 过滤:替代旧 list_active + 内存 retain 全表扫(任务量增长后
/// N×M 热点,每页都拉全表进内存再丢)。list_tasks 在有 project_id 时优先走本方法。
pub async fn list_active_by_project(&self, project_id: &str) -> Result<Vec<TaskRecord>> {
let conn = self.conn.clone();
let pid = project_id.to_owned();
tokio::task::spawn_blocking(move || {
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, output_json, idea_id, queue, parent_id, content_json, created_at, updated_at FROM tasks WHERE deleted_at IS NULL AND project_id = ?1 ORDER BY created_at DESC")
.map_err(storage_err)?;
let rows = stmt
.query_map(params![pid], |row| task_from_row(row))
.map_err(storage_err)?;
let mut results = Vec::new();
for r in rows {
results.push(r.map_err(storage_err)?);
}
Ok(results)
})
.await
.map_err(storage_err)?
}
/// 动态条件列出未删除任务(F-260621-02 查询维度补全)。
///
/// 复用 KnowledgeRepo::search 的「动态 WHERE + 参数绑定」模式,但用累积式条件收集
/// (Vec<String> WHERE 子句 + Vec<rusqlite::Value> 参数)替代 if-let 二分支——
/// TaskQuery 有 4 个过滤维度(project_id/status/priority/assignee/keyword),2^n 分支不可行,
/// 累积式天然支持任意维度组合,且每个 if-let 分支只 push 子句+参数,新增维度零样板。
///
/// - 过滤维度:project_id / status / priority / assignee(精确等值)+ keyword(title/description LIKE)
/// - keyword 拼成 `(title LIKE ?N OR description LIKE ?M)`,pattern = `%kw%`(对齐知识库 search)
/// - order_by 白名单(validate_order_by 防 SQL 注入),默认 created_at,恒 DESC(与 list_active 一致)
/// - limit/offset 钳制(limit ≤500 防滥用,对齐 conversation_repo::list_recent 的 limit≤200 思路)
/// - deleted_at IS NULL 恒带(回收站任务不进结果,语义同 list_active,不可被 query 关闭)
/// - 显式列出全部 18 列(不 SELECT deleted_at:TaskRecord 不带该字段,取了 from_row 报未知列)
///
/// 空 query(全 None)→ 等价 list_active(全量未删,created_at DESC),向后兼容。
/// status 值合法性由上层 list_tasks 命令(TaskStatus::is_valid)兜底,本层不过滤值集
/// (非法 status 在 DB 无匹配行,返回空 Vec,无害)。
pub async fn list_by_query(&self, query: &TaskQuery) -> Result<Vec<TaskRecord>> {
let conn = self.conn.clone();
// order_by 白名单校验在闭包外做(提前 fail-fast,非法值不进 DB 层)。列名不可
// 参数化,只能白名单防注入(对齐 impl_repo! 宏 validate_column_name 思路)。
// order_col 转 String 拥有所有权:避免 &str 借用 query 跨 spawn_blocking 'static
// 闭包(E0521 borrowed data escapes)。
let order_col = query.order_by.as_deref().unwrap_or("created_at");
validate_order_by(order_col)?;
let order_col = order_col.to_string();
// limit 钳制上限 500 防滥用(对齐 conversation_repo::list_recent 的 limit.min(200)
// 思路,任务场景放宽;None → 不拼 LIMIT = 全量,语义同 list_active)。
let limit_i: Option<i64> = query.limit.map(|l| (l.min(500)) as i64);
let offset_i: i64 = query.offset.unwrap_or(0) as i64;
// 拷贝 query 字段进闭包('static 生命周期,spawn_blocking 要求)
let project_id = query.project_id.clone();
let status = query.status.clone();
let priority = query.priority;
let assignee = query.assignee.clone();
let keyword = query.keyword.clone();
let queue = query.queue.clone();
let parent_id = query.parent_id.clone();
tokio::task::spawn_blocking(move || {
let guard = conn.blocking_lock();
// ── 累积 WHERE 子句 + 收集参数(按出现顺序绑定占位符 ?N,序号 = params_vec.len()+1)──
// 占位符序号必须按「实际参数位置」(params_vec.len()+1)而非「子句数」(where_clauses.len()+1)
// 编号:deleted_at IS NULL 是常量条件无占位符却占 where_clauses[0],用子句数编号会让首个
// 真参数拿到 ?2 而 params_vec 只有 1 个元素 → rusqlite "needed 2, got 1"。
// 用 params_vec.len()+1 保证 ?N 与参数位置严格对齐(N = 参数序号)。
// deleted_at IS NULL 恒带(常量条件无占位符),回收站任务不进结果(语义同 list_active)。
let mut where_clauses: Vec<String> = vec!["deleted_at IS NULL".to_string()];
let mut params_vec: Vec<Box<dyn rusqlite::ToSql>> = Vec::new();
if let Some(pid) = &project_id {
where_clauses.push(format!("project_id = ?{}", params_vec.len() + 1));
params_vec.push(Box::new(pid.clone()));
}
if let Some(s) = &status {
where_clauses.push(format!("status = ?{}", params_vec.len() + 1));
params_vec.push(Box::new(s.clone()));
}
if let Some(p) = priority {
where_clauses.push(format!("priority = ?{}", params_vec.len() + 1));
params_vec.push(Box::new(p));
}
if let Some(a) = &assignee {
where_clauses.push(format!("assignee = ?{}", params_vec.len() + 1));
params_vec.push(Box::new(a.clone()));
}
// queue(知识图谱 Phase 1 V29):管理池等值过滤,看板视图按池分列数据源
if let Some(q) = &queue {
where_clauses.push(format!("queue = ?{}", params_vec.len() + 1));
params_vec.push(Box::new(q.clone()));
}
// parent_id(知识图谱 Phase 1 V29):父任务等值过滤(查某父的子任务)
if let Some(pid) = &parent_id {
where_clauses.push(format!("parent_id = ?{}", params_vec.len() + 1));
params_vec.push(Box::new(pid.clone()));
}
// keyword: title/description LIKE %kw%(P2,对齐知识库 search 的 LIKE 模式)
if let Some(kw) = &keyword {
let escaped = kw.replace('%', "\\%").replace('_', "\\_");
let pat = format!("%{escaped}%");
let p1 = params_vec.len() + 1;
let p2 = p1 + 1;
where_clauses.push(format!("(title LIKE ?{p1} OR description LIKE ?{p2}) ESCAPE '\\'"));
params_vec.push(Box::new(pat.clone()));
params_vec.push(Box::new(pat));
}
// LIMIT/OFFSET:limit 为 None → 不拼(全量);有 limit 时 offset 紧跟其后。
let where_param_count = params_vec.len();
let limit_sql_bound = match limit_i {
Some(_) => format!(
" LIMIT ?{} OFFSET ?{}",
where_param_count + 1,
where_param_count + 2
),
None => String::new(),
};
// 拼 SQL:?N 占位符序号与 params_vec 顺序严格对应(累积时按 +1 递增保证)。
// 显式列出全部 18 列(含 V29 queue/parent_id/content_json,不 SELECT deleted_at:
// TaskRecord 不带该字段,取了 from_row 会因未知列报错)。
let sql = format!(
"SELECT 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 \
FROM tasks WHERE {} ORDER BY {} DESC{}",
where_clauses.join(" AND "),
order_col,
limit_sql_bound
);
let mut stmt = guard.prepare(&sql).map_err(storage_err)?;
// 组装参数引用数组(where 参数 + 可选 limit/offset)。
// limit/offset 也压入 params_vec 收口:借用引用需指向同一生命周期存活处,
// 收口到 params_vec 后再统一取引用,避免局部 l 生命周期不足(E0597,
// 对齐 idea_repo.rs list_by_query 同名模式)。
if let Some(l) = limit_i {
params_vec.push(Box::new(l));
params_vec.push(Box::new(offset_i));
}
let param_refs: Vec<&dyn rusqlite::ToSql> =
params_vec.iter().map(|p| p.as_ref()).collect();
let rows = stmt
.query_map(param_refs.as_slice(), |row| task_from_row(row))
.map_err(storage_err)?;
let mut results = Vec::new();
for r in rows {
results.push(r.map_err(storage_err)?);
}
Ok(results)
})
.await
.map_err(storage_err)?
}
/// 查某父任务的全部子任务(parent_id = ?,deleted_at IS NULL),按创建时间升序。
///
/// 知识图谱 Phase 1 V29(对标设计 §2.1 父任务聚合规则):父任务=容器模型,status
/// 不走状态机,由子任务聚合计算。本方法取子任务列表供聚合规则消费。
///
/// 注:嵌套深度限制 1 级(无孙任务)由 IPC 层校验不进 DB 约束,本方法只做等值查询。
pub async fn get_children(&self, parent_id: &str) -> Result<Vec<TaskRecord>> {
let conn = self.conn.clone();
let pid = parent_id.to_owned();
tokio::task::spawn_blocking(move || {
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, output_json, idea_id, queue, parent_id, content_json, created_at, updated_at FROM tasks WHERE deleted_at IS NULL AND parent_id = ?1 ORDER BY created_at ASC")
.map_err(storage_err)?;
let rows = stmt
.query_map(params![pid], |row| task_from_row(row))
.map_err(storage_err)?;
let mut results = Vec::new();
for r in rows {
results.push(r.map_err(storage_err)?);
}
Ok(results)
})
.await
.map_err(storage_err)?
}
/// 父任务聚合:按 status 分组统计子任务计数(对标设计 §2.1 聚合规则)。
///
/// 父任务 status 由子任务聚合计算(不走状态机),聚合规则:
/// - 全 todo → 父 todo
/// - 任一 in_progress → 父 in_progress
/// - 任一 blocked → 父 blocked
/// - 全 done/cancelled → 父 done
///
/// 本方法返回 `Vec<(status, count)>`(SQL GROUP BY 一次查询,数据量小 ~50 无压力),
/// 聚合规则的具体判定由调用方(commands 层)实现 —— 本层只提供原始计数,不持有
/// 业务聚合逻辑(CRUD 层只懂表/列语义,对标 B-260616-16 跨表校验下沉思路)。
pub async fn count_children_by_status(
&self,
parent_id: &str,
) -> Result<Vec<(String, i64)>> {
let conn = self.conn.clone();
let pid = parent_id.to_owned();
tokio::task::spawn_blocking(move || {
let guard = conn.blocking_lock();
let mut stmt = guard
.prepare(
"SELECT status, COUNT(*) AS cnt FROM tasks \
WHERE deleted_at IS NULL AND parent_id = ?1 \
GROUP BY status",
)
.map_err(storage_err)?;
let rows = stmt
.query_map(params![pid], |row| {
Ok((row.get::<_, String>(0)?, row.get::<_, i64>(1)?))
})
.map_err(storage_err)?;
let mut results = Vec::new();
for r in rows {
results.push(r.map_err(storage_err)?);
}
Ok(results)
})
.await
.map_err(storage_err)?
}
/// 父任务聚合专用 status 写入(知识图谱 Phase 1 V29,对标设计 §2.1 D3 父任务=容器模型)。
///
/// **这是父任务 status 的唯一写入路径**,绕过 D-260616-04 status 收口(通用 update_field
/// 白名单不含 status,所有叶子任务 status 走 advance_status_atomic 状态机)。父任务 status
/// **不走状态机**(容器模型,由子任务聚合计算),故需专用写入路径。
///
/// 防护:
/// - 方法名 `set_status_for_aggregation` 显式表明语义,非通用 setter,防误用。
/// - 调用方(commands::task::recompute_parent_status)负责聚合规则计算,本方法只落库。
/// - 不动 review_rounds(父任务不执行工作流,无 review 退回语义)。
///
/// 返回是否命中(父任务不存在/已删 → false)。
pub async fn set_status_for_aggregation(&self, id: &str, new_status: &str) -> Result<bool> {
let conn = self.conn.clone();
let id = id.to_owned();
let new_status = new_status.to_owned();
let now = now_millis_str();
tokio::task::spawn_blocking(move || {
let guard = conn.blocking_lock();
let affected = guard
.execute(
"UPDATE tasks SET status = ?1, updated_at = ?2 \
WHERE id = ?3 AND deleted_at IS NULL",
params![new_status, now, id],
)
.map_err(storage_err)?;
Ok(affected > 0)
})
.await
.map_err(storage_err)?
}
/// 列出回收站(deleted_at IS NOT NULL),按更新时间(≈删除时间)降序。对标 ProjectRepo::list_deleted。
///
/// 注:按项目列活跃任务走 list_active_by_project(SQL 下推 project_id),
/// 无 pid 时 fallback list_active。
pub async fn list_deleted(&self) -> Result<Vec<TaskRecord>> {
let conn = self.conn.clone();
tokio::task::spawn_blocking(move || {
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, output_json, idea_id, queue, parent_id, content_json, created_at, updated_at FROM tasks WHERE deleted_at IS NOT NULL ORDER BY updated_at DESC")
.map_err(storage_err)?;
let rows = stmt
.query_map([], |row| task_from_row(row))
.map_err(storage_err)?;
let mut results = Vec::new();
for r in rows {
results.push(r.map_err(storage_err)?);
}
Ok(results)
})
.await
.map_err(storage_err)?
}
}
// ============================================================
// 单元测试 — 知识图谱 Phase 1:queue/parent_id 筛选 + get_children(内存 DB)
// ============================================================
#[cfg(test)]
mod tests {
use super::*;
use crate::crud::ProjectRepo;
use crate::models::ProjectRecord;
use df_types::types::{ProjectStatus, TaskStatus};
/// 构造一条 TaskRecord fixture(queue/parent_id/status 可定制,V29 新维度 + 聚合测试用 status)。
fn trec(id: &str, queue: &str, parent_id: Option<&str>) -> TaskRecord {
trec_full(id, queue, parent_id, TaskStatus::Todo)
}
/// 全参 fixture(聚合测试需自定义 status 时用)。
fn trec_full(id: &str, queue: &str, parent_id: Option<&str>, status: TaskStatus) -> TaskRecord {
TaskRecord {
id: id.to_string(),
project_id: "proj-1".to_string(),
title: format!("task-{id}"),
description: String::new(),
status,
priority: 1,
branch_name: None,
assignee: None,
workflow_def_id: None,
base_branch: None,
review_rounds: 0,
output_json: None,
idea_id: None,
queue: queue.to_string(),
parent_id: parent_id.map(|s| s.to_string()),
content_json: None,
created_at: "1700000000000".to_string(),
updated_at: "1700000000000".to_string(),
}
}
/// 构造内存 DB + 占位 project(满足 tasks.project_id FK,PRAGMA foreign_keys=ON)。
async fn setup() -> TaskRepo {
let db = Database::open_in_memory().await.expect("open_in_memory");
let project_repo = ProjectRepo::new(&db);
project_repo
.insert(ProjectRecord {
id: "proj-1".to_string(),
name: "proj-1".to_string(),
description: String::new(),
status: ProjectStatus::Planning,
idea_id: None,
path: None,
stack: None,
created_at: "1700000000000".to_string(),
updated_at: "1700000000000".to_string(),
})
.await
.unwrap();
TaskRepo::new(&db)
}
#[tokio::test]
async fn list_by_query_queue_filter() {
let repo = setup().await;
repo.insert(trec("t1", "backlog", None)).await.unwrap();
repo.insert(trec("t2", "todo", None)).await.unwrap();
repo.insert(trec("t3", "todo", None)).await.unwrap();
repo.insert(trec("t4", "done", None)).await.unwrap();
// queue=todo → 只返回 t2/t3
let q = TaskQuery {
queue: Some("todo".to_string()),
..Default::default()
};
let res = repo.list_by_query(&q).await.unwrap();
let ids: Vec<_> = res.iter().map(|r| r.id.as_str()).collect();
assert_eq!(ids.len(), 2);
assert!(ids.contains(&"t2"));
assert!(ids.contains(&"t3"));
// queue=backlog → 只返回 t1
let q = TaskQuery {
queue: Some("backlog".to_string()),
..Default::default()
};
let res = repo.list_by_query(&q).await.unwrap();
let ids: Vec<_> = res.iter().map(|r| r.id.as_str()).collect();
assert_eq!(ids, vec!["t1"]);
}
#[tokio::test]
async fn list_by_query_parent_id_filter() {
let repo = setup().await;
// parent 父任务(叶子,parent_id=None)+ 3 子任务(parent_id=parent)
repo.insert(trec("parent", "todo", None)).await.unwrap();
repo.insert(trec("c1", "todo", Some("parent"))).await.unwrap();
repo.insert(trec("c2", "todo", Some("parent"))).await.unwrap();
repo.insert(trec("orphan", "todo", None)).await.unwrap();
// parent_id=parent → 只返回 c1/c2(父任务自身不匹配)
let q = TaskQuery {
parent_id: Some("parent".to_string()),
..Default::default()
};
let res = repo.list_by_query(&q).await.unwrap();
let ids: Vec<_> = res.iter().map(|r| r.id.as_str()).collect();
assert_eq!(ids.len(), 2);
assert!(ids.contains(&"c1"));
assert!(ids.contains(&"c2"));
}
#[tokio::test]
async fn list_by_query_queue_and_parent_id_combined() {
let repo = setup().await;
repo.insert(trec("parent", "todo", None)).await.unwrap();
repo.insert(trec("c1", "todo", Some("parent"))).await.unwrap();
repo.insert(trec("c2", "backlog", Some("parent"))).await.unwrap();
// queue=todo AND parent_id=parent → 只 c1(c2 是 backlog)
let q = TaskQuery {
queue: Some("todo".to_string()),
parent_id: Some("parent".to_string()),
..Default::default()
};
let res = repo.list_by_query(&q).await.unwrap();
let ids: Vec<_> = res.iter().map(|r| r.id.as_str()).collect();
assert_eq!(ids, vec!["c1"]);
}
#[tokio::test]
async fn list_by_query_empty_returns_all_active() {
// 空 query(全 None)→ 等价 list_active 全量未删,向后兼容
let repo = setup().await;
repo.insert(trec("t1", "todo", None)).await.unwrap();
repo.insert(trec("t2", "backlog", None)).await.unwrap();
let res = repo.list_by_query(&TaskQuery::default()).await.unwrap();
assert_eq!(res.len(), 2);
}
#[tokio::test]
async fn get_children_returns_only_direct_children() {
let repo = setup().await;
repo.insert(trec("parent", "todo", None)).await.unwrap();
repo.insert(trec("c1", "todo", Some("parent"))).await.unwrap();
repo.insert(trec("c2", "in_progress", Some("parent"))).await.unwrap();
repo.insert(trec("orphan", "todo", None)).await.unwrap();
let children = repo.get_children("parent").await.unwrap();
let ids: Vec<_> = children.iter().map(|r| r.id.as_str()).collect();
assert_eq!(ids.len(), 2);
assert!(ids.contains(&"c1"));
assert!(ids.contains(&"c2"));
}
#[tokio::test]
async fn get_children_empty_when_no_children() {
let repo = setup().await;
repo.insert(trec("parent", "todo", None)).await.unwrap();
let children = repo.get_children("parent").await.unwrap();
assert!(children.is_empty());
}
#[tokio::test]
async fn count_children_by_status_groups() {
let repo = setup().await;
repo.insert(trec("parent", "todo", None)).await.unwrap();
// 4 个子任务:status 分布 todo×2 / in_progress×1 / done×1(queue 统一 todo)
repo.insert(trec_full("c1", "todo", Some("parent"), "todo"))
.await
.unwrap();
repo.insert(trec_full("c2", "todo", Some("parent"), "todo"))
.await
.unwrap();
repo.insert(trec_full("c3", "todo", Some("parent"), "in_progress"))
.await
.unwrap();
repo.insert(trec_full("c4", "todo", Some("parent"), "done"))
.await
.unwrap();
let counts = repo.count_children_by_status("parent").await.unwrap();
// 转 map 便于断言(顺序由 GROUP BY 决定,不依赖)。父任务自身不计入(parent_id 非自身)。
let map: std::collections::HashMap<String, i64> = counts.into_iter().collect();
assert_eq!(map.get("todo"), Some(&2));
assert_eq!(map.get("in_progress"), Some(&1));
assert_eq!(map.get("done"), Some(&1));
assert!(!map.contains_key("cancelled"));
}
#[tokio::test]
async fn count_children_by_status_empty_when_no_children() {
let repo = setup().await;
repo.insert(trec("parent", "todo", None)).await.unwrap();
let counts = repo.count_children_by_status("parent").await.unwrap();
assert!(counts.is_empty());
}
#[tokio::test]
async fn get_children_excludes_soft_deleted() {
// 软删子任务不进 get_children 结果(deleted_at IS NULL 过滤)
let repo = setup().await;
repo.insert(trec("parent", "todo", None)).await.unwrap();
repo.insert(trec("c1", "todo", Some("parent"))).await.unwrap();
repo.insert(trec("c2", "todo", Some("parent"))).await.unwrap();
repo.soft_delete("c2").await.unwrap();
let children = repo.get_children("parent").await.unwrap();
let ids: Vec<_> = children.iter().map(|r| r.id.as_str()).collect();
assert_eq!(ids, vec!["c1"], "软删子任务应被过滤");
}
// ============================================================
// 父聚合专用 status 写入:set_status_for_aggregation(知识图谱 Phase 1 V29)
// 父任务=容器模型,status 不走状态机,由子任务聚合计算后经此方法落库。
// 锁定:① 写入命中 + status 变更;② 不动 review_rounds;③ 软删任务返回 false。
// ============================================================
#[tokio::test]
async fn set_status_for_aggregation_writes_status() {
let repo = setup().await;
// 父任务初始 todo(queue=todo, parent_id=None 容器模型)
repo.insert(trec_full("parent", "todo", None, TaskStatus::Todo))
.await
.unwrap();
let ok = repo.set_status_for_aggregation("parent", "in_progress");
assert!(ok, "应命中写入");
let after = repo.get_by_id("parent").await.unwrap().unwrap();
assert_eq!(after.status.as_str(), "in_progress", "status 应被聚合写入更新");
assert_eq!(after.review_rounds, 0, "父聚合写入不动 review_rounds");
}
#[tokio::test]
async fn set_status_for_aggregation_skips_soft_deleted() {
// 软删任务(回收站)不进聚合写入(WHERE deleted_at IS NULL),返回 false
let repo = setup().await;
repo.insert(trec_full("parent", "todo", None, TaskStatus::Todo))
.await
.unwrap();
repo.soft_delete("parent").await.unwrap();
let ok = repo.set_status_for_aggregation("parent", "done").await.unwrap();
assert!(!ok, "软删任务不应被聚合写入命中");
}
#[tokio::test]
async fn set_status_for_aggregation_nonexistent_returns_false() {
let repo = setup().await;
let ok = repo.set_status_for_aggregation("ghost", "done").await.unwrap();
assert!(!ok, "不存在的任务应返回 false");
}
}