新增: 多Agent数据层(V36迁移+3新表+Repo CRUD+7个测试全绿)
- V36: ai_plans/ai_subtasks/ai_conflicts 3 新表 + ai_messages/ai_tool_executions 加 subtask_id - models: PlanRecord/SubTaskRecord/ConflictRecord 结构体 - PlanRepo/SubTaskRepo/ConflictRepo CRUD(insert/get/list_by_*/update_status/resolve) - 7个 Repo 单元测试全部通过(内存 SQLite + 全迁移) - 顺带修复 task_repo.rs 预存测试编译错误(trec_full 参数类型+缺失 await)
This commit is contained in:
2
Batch.md
2
Batch.md
@@ -437,6 +437,8 @@
|
|||||||
|
|
||||||
### 数据层 + Git worktree 隔离 + 并行执行
|
### 数据层 + Git worktree 隔离 + 并行执行
|
||||||
|
|
||||||
|
- **提交**: `当前待提交`
|
||||||
|
|
||||||
| # | 任务 | 文件 |
|
| # | 任务 | 文件 |
|
||||||
|---|------|------|
|
|---|------|------|
|
||||||
| 1 | V36 迁移(ai_plans/ai_subtasks/ai_conflicts 3 新表 + ai_messages/ai_tool_executions 加 subtask_id + subtasks.branch + conflicts.conflict_type) | migrations.rs |
|
| 1 | V36 迁移(ai_plans/ai_subtasks/ai_conflicts 3 新表 + ai_messages/ai_tool_executions 加 subtask_id + subtasks.branch + conflicts.conflict_type) | migrations.rs |
|
||||||
|
|||||||
@@ -17,6 +17,9 @@
|
|||||||
//! re-export(`pub use ...::*`)保持 `df_storage::crud::XxxRepo` /
|
//! re-export(`pub use ...::*`)保持 `df_storage::crud::XxxRepo` /
|
||||||
//! `df_storage::crud::is_allowed_column` 路径不变,**调用方零改动**。
|
//! `df_storage::crud::is_allowed_column` 路径不变,**调用方零改动**。
|
||||||
|
|
||||||
|
mod module_dependency_repo;
|
||||||
|
mod plan_repo;
|
||||||
|
|
||||||
mod conversation_repo;
|
mod conversation_repo;
|
||||||
mod idea_eval_repo;
|
mod idea_eval_repo;
|
||||||
mod idea_repo;
|
mod idea_repo;
|
||||||
@@ -28,7 +31,6 @@ mod project_service_repo;
|
|||||||
mod settings;
|
mod settings;
|
||||||
mod task_link_repo;
|
mod task_link_repo;
|
||||||
mod task_repo;
|
mod task_repo;
|
||||||
mod module_dependency_repo;
|
|
||||||
|
|
||||||
pub use conversation_repo::*;
|
pub use conversation_repo::*;
|
||||||
pub use idea_eval_repo::*;
|
pub use idea_eval_repo::*;
|
||||||
@@ -42,6 +44,7 @@ pub use settings::*;
|
|||||||
pub use task_link_repo::*;
|
pub use task_link_repo::*;
|
||||||
pub use task_repo::*;
|
pub use task_repo::*;
|
||||||
pub use module_dependency_repo::*;
|
pub use module_dependency_repo::*;
|
||||||
|
pub use plan_repo::*;
|
||||||
|
|
||||||
// ============================================================
|
// ============================================================
|
||||||
// 辅助宏 — 消除多个 Repo 的重复样板
|
// 辅助宏 — 消除多个 Repo 的重复样板
|
||||||
|
|||||||
579
crates/df-storage/src/crud/plan_repo.rs
Normal file
579
crates/df-storage/src/crud/plan_repo.rs
Normal file
@@ -0,0 +1,579 @@
|
|||||||
|
//! 多 Agent 并行执行 Repo — ai_plans / ai_subtasks / ai_conflicts 表 CRUD(V36)
|
||||||
|
//!
|
||||||
|
//! 设计依据:docs/02-架构设计/专项设计/多Agent并行执行与仲裁合并设计-2026-07-01.md §二
|
||||||
|
|
||||||
|
use std::sync::Arc;
|
||||||
|
|
||||||
|
use rusqlite::{params, OptionalExtension, Row};
|
||||||
|
use tokio::sync::Mutex;
|
||||||
|
|
||||||
|
use crate::db::Database;
|
||||||
|
use crate::models::{ConflictRecord, PlanRecord, SubTaskRecord};
|
||||||
|
|
||||||
|
use super::storage_err;
|
||||||
|
|
||||||
|
type Result<T> = std::result::Result<T, df_types::error::Error>;
|
||||||
|
|
||||||
|
fn plan_from_row(row: &Row<'_>) -> rusqlite::Result<PlanRecord> {
|
||||||
|
Ok(PlanRecord {
|
||||||
|
id: row.get("id")?,
|
||||||
|
conversation_id: row.get("conversation_id")?,
|
||||||
|
user_message_id: row.get("user_message_id")?,
|
||||||
|
status: row.get("status")?,
|
||||||
|
subtask_count: row.get("subtask_count")?,
|
||||||
|
created_at: row.get("created_at")?,
|
||||||
|
completed_at: row.get("completed_at")?,
|
||||||
|
})
|
||||||
|
}
|
||||||
|
|
||||||
|
fn subtask_from_row(row: &Row<'_>) -> rusqlite::Result<SubTaskRecord> {
|
||||||
|
Ok(SubTaskRecord {
|
||||||
|
id: row.get("id")?,
|
||||||
|
plan_id: row.get("plan_id")?,
|
||||||
|
persona_id: row.get("persona_id")?,
|
||||||
|
intent: row.get("intent")?,
|
||||||
|
status: row.get("status")?,
|
||||||
|
layer: row.get("layer")?,
|
||||||
|
deps: row.get("deps")?,
|
||||||
|
branch: row.get("branch")?,
|
||||||
|
created_at: row.get("created_at")?,
|
||||||
|
completed_at: row.get("completed_at")?,
|
||||||
|
})
|
||||||
|
}
|
||||||
|
|
||||||
|
fn conflict_from_row(row: &Row<'_>) -> rusqlite::Result<ConflictRecord> {
|
||||||
|
Ok(ConflictRecord {
|
||||||
|
id: row.get("id")?,
|
||||||
|
plan_id: row.get("plan_id")?,
|
||||||
|
file_path: row.get("file_path")?,
|
||||||
|
conflict_type: row.get("conflict_type")?,
|
||||||
|
subtask_a: row.get("subtask_a")?,
|
||||||
|
subtask_b: row.get("subtask_b")?,
|
||||||
|
diff_a: row.get("diff_a")?,
|
||||||
|
diff_b: row.get("diff_b")?,
|
||||||
|
resolution: row.get("resolution")?,
|
||||||
|
resolved_by: row.get("resolved_by")?,
|
||||||
|
created_at: row.get("created_at")?,
|
||||||
|
resolved_at: row.get("resolved_at")?,
|
||||||
|
})
|
||||||
|
}
|
||||||
|
|
||||||
|
// ============================================================
|
||||||
|
// PlanRepo
|
||||||
|
// ============================================================
|
||||||
|
|
||||||
|
pub struct PlanRepo {
|
||||||
|
conn: Arc<Mutex<rusqlite::Connection>>,
|
||||||
|
}
|
||||||
|
|
||||||
|
impl PlanRepo {
|
||||||
|
pub fn new(db: &Database) -> Self {
|
||||||
|
Self { conn: db.conn() }
|
||||||
|
}
|
||||||
|
|
||||||
|
pub async fn insert(&self, record: PlanRecord) -> Result<String> {
|
||||||
|
let conn = self.conn.clone();
|
||||||
|
tokio::task::spawn_blocking(move || {
|
||||||
|
let guard = conn.blocking_lock();
|
||||||
|
guard
|
||||||
|
.execute(
|
||||||
|
"INSERT INTO ai_plans \
|
||||||
|
(id, conversation_id, user_message_id, status, subtask_count, created_at, completed_at) \
|
||||||
|
VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7)",
|
||||||
|
params![
|
||||||
|
record.id, record.conversation_id, record.user_message_id,
|
||||||
|
record.status, record.subtask_count, record.created_at, record.completed_at,
|
||||||
|
],
|
||||||
|
)
|
||||||
|
.map_err(storage_err)?;
|
||||||
|
Ok(record.id)
|
||||||
|
})
|
||||||
|
.await
|
||||||
|
.map_err(storage_err)?
|
||||||
|
}
|
||||||
|
|
||||||
|
pub async fn get(&self, id: &str) -> Result<Option<PlanRecord>> {
|
||||||
|
let conn = self.conn.clone();
|
||||||
|
let id = id.to_owned();
|
||||||
|
tokio::task::spawn_blocking(move || {
|
||||||
|
let guard = conn.blocking_lock();
|
||||||
|
let row = guard
|
||||||
|
.query_row("SELECT * FROM ai_plans WHERE id = ?1", params![id], plan_from_row)
|
||||||
|
.optional()
|
||||||
|
.map_err(storage_err)?;
|
||||||
|
Ok(row)
|
||||||
|
})
|
||||||
|
.await
|
||||||
|
.map_err(storage_err)?
|
||||||
|
}
|
||||||
|
|
||||||
|
pub async fn update_status(
|
||||||
|
&self,
|
||||||
|
id: &str,
|
||||||
|
status: &str,
|
||||||
|
completed_at: Option<&str>,
|
||||||
|
) -> Result<bool> {
|
||||||
|
let conn = self.conn.clone();
|
||||||
|
let id = id.to_owned();
|
||||||
|
let status = status.to_owned();
|
||||||
|
let completed_at = completed_at.map(|s| s.to_owned());
|
||||||
|
tokio::task::spawn_blocking(move || {
|
||||||
|
let guard = conn.blocking_lock();
|
||||||
|
let affected = guard
|
||||||
|
.execute(
|
||||||
|
"UPDATE ai_plans SET status = ?1, completed_at = ?2 WHERE id = ?3",
|
||||||
|
params![status, completed_at, id],
|
||||||
|
)
|
||||||
|
.map_err(storage_err)?;
|
||||||
|
Ok(affected > 0)
|
||||||
|
})
|
||||||
|
.await
|
||||||
|
.map_err(storage_err)?
|
||||||
|
}
|
||||||
|
|
||||||
|
pub async fn list_by_conversation(&self, conv_id: &str) -> Result<Vec<PlanRecord>> {
|
||||||
|
let conn = self.conn.clone();
|
||||||
|
let conv_id = conv_id.to_owned();
|
||||||
|
tokio::task::spawn_blocking(move || {
|
||||||
|
let guard = conn.blocking_lock();
|
||||||
|
let mut stmt = guard
|
||||||
|
.prepare("SELECT * FROM ai_plans WHERE conversation_id = ?1 ORDER BY created_at DESC")
|
||||||
|
.map_err(storage_err)?;
|
||||||
|
let rows = stmt.query_map(params![conv_id], plan_from_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)?
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// ============================================================
|
||||||
|
// SubTaskRepo
|
||||||
|
// ============================================================
|
||||||
|
|
||||||
|
pub struct SubTaskRepo {
|
||||||
|
conn: Arc<Mutex<rusqlite::Connection>>,
|
||||||
|
}
|
||||||
|
|
||||||
|
impl SubTaskRepo {
|
||||||
|
pub fn new(db: &Database) -> Self {
|
||||||
|
Self { conn: db.conn() }
|
||||||
|
}
|
||||||
|
|
||||||
|
pub async fn insert(&self, record: SubTaskRecord) -> Result<String> {
|
||||||
|
let conn = self.conn.clone();
|
||||||
|
tokio::task::spawn_blocking(move || {
|
||||||
|
let guard = conn.blocking_lock();
|
||||||
|
guard
|
||||||
|
.execute(
|
||||||
|
"INSERT INTO ai_subtasks \
|
||||||
|
(id, plan_id, persona_id, intent, status, layer, deps, branch, created_at, completed_at) \
|
||||||
|
VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, ?10)",
|
||||||
|
params![
|
||||||
|
record.id, record.plan_id, record.persona_id, record.intent,
|
||||||
|
record.status, record.layer, record.deps, record.branch,
|
||||||
|
record.created_at, record.completed_at,
|
||||||
|
],
|
||||||
|
)
|
||||||
|
.map_err(storage_err)?;
|
||||||
|
Ok(record.id)
|
||||||
|
})
|
||||||
|
.await
|
||||||
|
.map_err(storage_err)?
|
||||||
|
}
|
||||||
|
|
||||||
|
pub async fn list_by_plan(&self, plan_id: &str) -> Result<Vec<SubTaskRecord>> {
|
||||||
|
let conn = self.conn.clone();
|
||||||
|
let plan_id = plan_id.to_owned();
|
||||||
|
tokio::task::spawn_blocking(move || {
|
||||||
|
let guard = conn.blocking_lock();
|
||||||
|
let mut stmt = guard
|
||||||
|
.prepare("SELECT * FROM ai_subtasks WHERE plan_id = ?1 ORDER BY layer ASC, created_at ASC")
|
||||||
|
.map_err(storage_err)?;
|
||||||
|
let rows = stmt.query_map(params![plan_id], subtask_from_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)?
|
||||||
|
}
|
||||||
|
|
||||||
|
pub async fn update_status(
|
||||||
|
&self,
|
||||||
|
id: &str,
|
||||||
|
status: &str,
|
||||||
|
completed_at: Option<&str>,
|
||||||
|
) -> Result<bool> {
|
||||||
|
let conn = self.conn.clone();
|
||||||
|
let id = id.to_owned();
|
||||||
|
let status = status.to_owned();
|
||||||
|
let completed_at = completed_at.map(|s| s.to_owned());
|
||||||
|
tokio::task::spawn_blocking(move || {
|
||||||
|
let guard = conn.blocking_lock();
|
||||||
|
let affected = guard
|
||||||
|
.execute(
|
||||||
|
"UPDATE ai_subtasks SET status = ?1, completed_at = ?2 WHERE id = ?3",
|
||||||
|
params![status, completed_at, id],
|
||||||
|
)
|
||||||
|
.map_err(storage_err)?;
|
||||||
|
Ok(affected > 0)
|
||||||
|
})
|
||||||
|
.await
|
||||||
|
.map_err(storage_err)?
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// ============================================================
|
||||||
|
// ConflictRepo
|
||||||
|
// ============================================================
|
||||||
|
|
||||||
|
pub struct ConflictRepo {
|
||||||
|
conn: Arc<Mutex<rusqlite::Connection>>,
|
||||||
|
}
|
||||||
|
|
||||||
|
impl ConflictRepo {
|
||||||
|
pub fn new(db: &Database) -> Self {
|
||||||
|
Self { conn: db.conn() }
|
||||||
|
}
|
||||||
|
|
||||||
|
pub async fn insert(&self, record: ConflictRecord) -> Result<String> {
|
||||||
|
let conn = self.conn.clone();
|
||||||
|
tokio::task::spawn_blocking(move || {
|
||||||
|
let guard = conn.blocking_lock();
|
||||||
|
guard
|
||||||
|
.execute(
|
||||||
|
"INSERT INTO ai_conflicts \
|
||||||
|
(id, plan_id, file_path, conflict_type, subtask_a, subtask_b, diff_a, diff_b, \
|
||||||
|
resolution, resolved_by, created_at, resolved_at) \
|
||||||
|
VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, ?10, ?11, ?12)",
|
||||||
|
params![
|
||||||
|
record.id, record.plan_id, record.file_path, record.conflict_type,
|
||||||
|
record.subtask_a, record.subtask_b, record.diff_a, record.diff_b,
|
||||||
|
record.resolution, record.resolved_by, record.created_at, record.resolved_at,
|
||||||
|
],
|
||||||
|
)
|
||||||
|
.map_err(storage_err)?;
|
||||||
|
Ok(record.id)
|
||||||
|
})
|
||||||
|
.await
|
||||||
|
.map_err(storage_err)?
|
||||||
|
}
|
||||||
|
|
||||||
|
pub async fn list_pending(&self, plan_id: &str) -> Result<Vec<ConflictRecord>> {
|
||||||
|
let conn = self.conn.clone();
|
||||||
|
let plan_id = plan_id.to_owned();
|
||||||
|
tokio::task::spawn_blocking(move || {
|
||||||
|
let guard = conn.blocking_lock();
|
||||||
|
let mut stmt = guard
|
||||||
|
.prepare("SELECT * FROM ai_conflicts WHERE plan_id = ?1 AND resolution = 'pending' ORDER BY created_at ASC")
|
||||||
|
.map_err(storage_err)?;
|
||||||
|
let rows = stmt.query_map(params![plan_id], conflict_from_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)?
|
||||||
|
}
|
||||||
|
|
||||||
|
pub async fn resolve(
|
||||||
|
&self,
|
||||||
|
id: &str,
|
||||||
|
resolution: &str,
|
||||||
|
resolved_by: &str,
|
||||||
|
resolved_at: &str,
|
||||||
|
) -> Result<bool> {
|
||||||
|
let conn = self.conn.clone();
|
||||||
|
let id = id.to_owned();
|
||||||
|
let resolution = resolution.to_owned();
|
||||||
|
let resolved_by = resolved_by.to_owned();
|
||||||
|
let resolved_at = resolved_at.to_owned();
|
||||||
|
tokio::task::spawn_blocking(move || {
|
||||||
|
let guard = conn.blocking_lock();
|
||||||
|
let affected = guard
|
||||||
|
.execute(
|
||||||
|
"UPDATE ai_conflicts SET resolution = ?1, resolved_by = ?2, resolved_at = ?3 WHERE id = ?4",
|
||||||
|
params![resolution, resolved_by, resolved_at, id],
|
||||||
|
)
|
||||||
|
.map_err(storage_err)?;
|
||||||
|
Ok(affected > 0)
|
||||||
|
})
|
||||||
|
.await
|
||||||
|
.map_err(storage_err)?
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
#[cfg(test)]
|
||||||
|
mod tests {
|
||||||
|
use super::*;
|
||||||
|
use crate::db::Database;
|
||||||
|
|
||||||
|
async fn setup() -> Database {
|
||||||
|
let db = Database::open_in_memory().await.expect("open_in_memory");
|
||||||
|
let conn = db.conn();
|
||||||
|
tokio::task::spawn_blocking(move || {
|
||||||
|
crate::migrations::run(&conn.blocking_lock()).expect("migrations");
|
||||||
|
})
|
||||||
|
.await
|
||||||
|
.expect("spawn_blocking");
|
||||||
|
db
|
||||||
|
}
|
||||||
|
|
||||||
|
#[tokio::test]
|
||||||
|
async fn repo_01_plan_insert_and_get() {
|
||||||
|
let db = setup().await;
|
||||||
|
let repo = PlanRepo::new(&db);
|
||||||
|
let plan = PlanRecord {
|
||||||
|
id: "plan-001".into(),
|
||||||
|
conversation_id: "conv-001".into(),
|
||||||
|
user_message_id: Some("msg-001".into()),
|
||||||
|
status: "planning".into(),
|
||||||
|
subtask_count: 2,
|
||||||
|
created_at: "2026-07-01T00:00:00Z".into(),
|
||||||
|
completed_at: None,
|
||||||
|
};
|
||||||
|
let id = repo.insert(plan.clone()).await.unwrap();
|
||||||
|
assert_eq!(id, "plan-001");
|
||||||
|
let got = repo.get("plan-001").await.unwrap().unwrap();
|
||||||
|
assert_eq!(got.status, "planning");
|
||||||
|
assert_eq!(got.subtask_count, 2);
|
||||||
|
}
|
||||||
|
|
||||||
|
#[tokio::test]
|
||||||
|
async fn repo_02_plan_update_status() {
|
||||||
|
let db = setup().await;
|
||||||
|
let repo = PlanRepo::new(&db);
|
||||||
|
repo.insert(PlanRecord {
|
||||||
|
id: "plan-002".into(),
|
||||||
|
conversation_id: "conv-001".into(),
|
||||||
|
user_message_id: None,
|
||||||
|
status: "planning".into(),
|
||||||
|
subtask_count: 0,
|
||||||
|
created_at: "2026-07-01T00:00:00Z".into(),
|
||||||
|
completed_at: None,
|
||||||
|
})
|
||||||
|
.await
|
||||||
|
.unwrap();
|
||||||
|
assert!(repo.update_status("plan-002", "executing", None).await.unwrap());
|
||||||
|
let got = repo.get("plan-002").await.unwrap().unwrap();
|
||||||
|
assert_eq!(got.status, "executing");
|
||||||
|
}
|
||||||
|
|
||||||
|
#[tokio::test]
|
||||||
|
async fn repo_03_subtask_insert_and_list() {
|
||||||
|
let db = setup().await;
|
||||||
|
let plan_repo = PlanRepo::new(&db);
|
||||||
|
let st_repo = SubTaskRepo::new(&db);
|
||||||
|
plan_repo
|
||||||
|
.insert(PlanRecord {
|
||||||
|
id: "plan-003".into(),
|
||||||
|
conversation_id: "conv-001".into(),
|
||||||
|
user_message_id: None,
|
||||||
|
status: "planning".into(),
|
||||||
|
subtask_count: 2,
|
||||||
|
created_at: "2026-07-01T00:00:00Z".into(),
|
||||||
|
completed_at: None,
|
||||||
|
})
|
||||||
|
.await
|
||||||
|
.unwrap();
|
||||||
|
st_repo
|
||||||
|
.insert(SubTaskRecord {
|
||||||
|
id: "st-001".into(),
|
||||||
|
plan_id: "plan-003".into(),
|
||||||
|
persona_id: Some("coder".into()),
|
||||||
|
intent: "读取代码".into(),
|
||||||
|
status: "pending".into(),
|
||||||
|
layer: 0,
|
||||||
|
deps: None,
|
||||||
|
branch: Some("subtask/plan-003/st-001".into()),
|
||||||
|
created_at: "2026-07-01T00:00:00Z".into(),
|
||||||
|
completed_at: None,
|
||||||
|
})
|
||||||
|
.await
|
||||||
|
.unwrap();
|
||||||
|
st_repo
|
||||||
|
.insert(SubTaskRecord {
|
||||||
|
id: "st-002".into(),
|
||||||
|
plan_id: "plan-003".into(),
|
||||||
|
persona_id: Some("coder".into()),
|
||||||
|
intent: "修改代码".into(),
|
||||||
|
status: "pending".into(),
|
||||||
|
layer: 1,
|
||||||
|
deps: Some(r#"["st-001"]"#.into()),
|
||||||
|
branch: Some("subtask/plan-003/st-002".into()),
|
||||||
|
created_at: "2026-07-01T00:00:00Z".into(),
|
||||||
|
completed_at: None,
|
||||||
|
})
|
||||||
|
.await
|
||||||
|
.unwrap();
|
||||||
|
let list = st_repo.list_by_plan("plan-003").await.unwrap();
|
||||||
|
assert_eq!(list.len(), 2);
|
||||||
|
assert_eq!(list[0].layer, 0); // layer 排序
|
||||||
|
assert_eq!(list[1].layer, 1);
|
||||||
|
}
|
||||||
|
|
||||||
|
#[tokio::test]
|
||||||
|
async fn repo_04_subtask_update_status() {
|
||||||
|
let db = setup().await;
|
||||||
|
let plan_repo = PlanRepo::new(&db);
|
||||||
|
let st_repo = SubTaskRepo::new(&db);
|
||||||
|
plan_repo
|
||||||
|
.insert(PlanRecord {
|
||||||
|
id: "plan-004".into(),
|
||||||
|
conversation_id: "conv-001".into(),
|
||||||
|
user_message_id: None,
|
||||||
|
status: "planning".into(),
|
||||||
|
subtask_count: 1,
|
||||||
|
created_at: "2026-07-01T00:00:00Z".into(),
|
||||||
|
completed_at: None,
|
||||||
|
})
|
||||||
|
.await
|
||||||
|
.unwrap();
|
||||||
|
st_repo
|
||||||
|
.insert(SubTaskRecord {
|
||||||
|
id: "st-004".into(),
|
||||||
|
plan_id: "plan-004".into(),
|
||||||
|
persona_id: None,
|
||||||
|
intent: "test".into(),
|
||||||
|
status: "pending".into(),
|
||||||
|
layer: 0,
|
||||||
|
deps: None,
|
||||||
|
branch: None,
|
||||||
|
created_at: "2026-07-01T00:00:00Z".into(),
|
||||||
|
completed_at: None,
|
||||||
|
})
|
||||||
|
.await
|
||||||
|
.unwrap();
|
||||||
|
assert!(st_repo.update_status("st-004", "done", Some("2026-07-01T01:00:00Z")).await.unwrap());
|
||||||
|
let list = st_repo.list_by_plan("plan-004").await.unwrap();
|
||||||
|
assert_eq!(list[0].status, "done");
|
||||||
|
}
|
||||||
|
|
||||||
|
#[tokio::test]
|
||||||
|
async fn repo_05_conflict_insert_and_list_pending() {
|
||||||
|
let db = setup().await;
|
||||||
|
let plan_repo = PlanRepo::new(&db);
|
||||||
|
let conf_repo = ConflictRepo::new(&db);
|
||||||
|
plan_repo
|
||||||
|
.insert(PlanRecord {
|
||||||
|
id: "plan-005".into(),
|
||||||
|
conversation_id: "conv-001".into(),
|
||||||
|
user_message_id: None,
|
||||||
|
status: "merging".into(),
|
||||||
|
subtask_count: 2,
|
||||||
|
created_at: "2026-07-01T00:00:00Z".into(),
|
||||||
|
completed_at: None,
|
||||||
|
})
|
||||||
|
.await
|
||||||
|
.unwrap();
|
||||||
|
conf_repo
|
||||||
|
.insert(ConflictRecord {
|
||||||
|
id: "conf-001".into(),
|
||||||
|
plan_id: "plan-005".into(),
|
||||||
|
file_path: "src/main.rs".into(),
|
||||||
|
conflict_type: "file".into(),
|
||||||
|
subtask_a: Some("st-a".into()),
|
||||||
|
subtask_b: Some("st-b".into()),
|
||||||
|
diff_a: Some("-old\n+new_a".into()),
|
||||||
|
diff_b: Some("-old\n+new_b".into()),
|
||||||
|
resolution: "pending".into(),
|
||||||
|
resolved_by: None,
|
||||||
|
created_at: "2026-07-01T00:00:00Z".into(),
|
||||||
|
resolved_at: None,
|
||||||
|
})
|
||||||
|
.await
|
||||||
|
.unwrap();
|
||||||
|
let pending = conf_repo.list_pending("plan-005").await.unwrap();
|
||||||
|
assert_eq!(pending.len(), 1);
|
||||||
|
assert_eq!(pending[0].file_path, "src/main.rs");
|
||||||
|
}
|
||||||
|
|
||||||
|
#[tokio::test]
|
||||||
|
async fn repo_06_conflict_resolve() {
|
||||||
|
let db = setup().await;
|
||||||
|
let plan_repo = PlanRepo::new(&db);
|
||||||
|
let conf_repo = ConflictRepo::new(&db);
|
||||||
|
plan_repo
|
||||||
|
.insert(PlanRecord {
|
||||||
|
id: "plan-006".into(),
|
||||||
|
conversation_id: "conv-001".into(),
|
||||||
|
user_message_id: None,
|
||||||
|
status: "merging".into(),
|
||||||
|
subtask_count: 1,
|
||||||
|
created_at: "2026-07-01T00:00:00Z".into(),
|
||||||
|
completed_at: None,
|
||||||
|
})
|
||||||
|
.await
|
||||||
|
.unwrap();
|
||||||
|
conf_repo
|
||||||
|
.insert(ConflictRecord {
|
||||||
|
id: "conf-002".into(),
|
||||||
|
plan_id: "plan-006".into(),
|
||||||
|
file_path: "src/lib.rs".into(),
|
||||||
|
conflict_type: "semantic".into(),
|
||||||
|
subtask_a: None,
|
||||||
|
subtask_b: None,
|
||||||
|
diff_a: None,
|
||||||
|
diff_b: None,
|
||||||
|
resolution: "pending".into(),
|
||||||
|
resolved_by: None,
|
||||||
|
created_at: "2026-07-01T00:00:00Z".into(),
|
||||||
|
resolved_at: None,
|
||||||
|
})
|
||||||
|
.await
|
||||||
|
.unwrap();
|
||||||
|
assert!(
|
||||||
|
conf_repo
|
||||||
|
.resolve("conf-002", "merged", "reviewer", "2026-07-01T02:00:00Z")
|
||||||
|
.await
|
||||||
|
.unwrap()
|
||||||
|
);
|
||||||
|
let pending = conf_repo.list_pending("plan-006").await.unwrap();
|
||||||
|
assert_eq!(pending.len(), 0, "解决后 pending 列表应为空");
|
||||||
|
}
|
||||||
|
|
||||||
|
#[tokio::test]
|
||||||
|
async fn repo_07_subtask_branch_null_for_non_git() {
|
||||||
|
let db = setup().await;
|
||||||
|
let plan_repo = PlanRepo::new(&db);
|
||||||
|
let st_repo = SubTaskRepo::new(&db);
|
||||||
|
plan_repo
|
||||||
|
.insert(PlanRecord {
|
||||||
|
id: "plan-007".into(),
|
||||||
|
conversation_id: "conv-001".into(),
|
||||||
|
user_message_id: None,
|
||||||
|
status: "planning".into(),
|
||||||
|
subtask_count: 1,
|
||||||
|
created_at: "2026-07-01T00:00:00Z".into(),
|
||||||
|
completed_at: None,
|
||||||
|
})
|
||||||
|
.await
|
||||||
|
.unwrap();
|
||||||
|
st_repo
|
||||||
|
.insert(SubTaskRecord {
|
||||||
|
id: "st-007".into(),
|
||||||
|
plan_id: "plan-007".into(),
|
||||||
|
persona_id: None,
|
||||||
|
intent: "非 Git 工程".into(),
|
||||||
|
status: "pending".into(),
|
||||||
|
layer: 0,
|
||||||
|
deps: None,
|
||||||
|
branch: None, // 非 Git 工程
|
||||||
|
created_at: "2026-07-01T00:00:00Z".into(),
|
||||||
|
completed_at: None,
|
||||||
|
})
|
||||||
|
.await
|
||||||
|
.unwrap();
|
||||||
|
let list = st_repo.list_by_plan("plan-007").await.unwrap();
|
||||||
|
assert!(list[0].branch.is_none(), "非 Git 工程 branch 应为 NULL");
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -784,16 +784,16 @@ mod tests {
|
|||||||
let repo = setup().await;
|
let repo = setup().await;
|
||||||
repo.insert(trec("parent", "todo", None)).await.unwrap();
|
repo.insert(trec("parent", "todo", None)).await.unwrap();
|
||||||
// 4 个子任务:status 分布 todo×2 / in_progress×1 / done×1(queue 统一 todo)
|
// 4 个子任务:status 分布 todo×2 / in_progress×1 / done×1(queue 统一 todo)
|
||||||
repo.insert(trec_full("c1", "todo", Some("parent"), "todo"))
|
repo.insert(trec_full("c1", "todo", Some("parent"), TaskStatus::Todo))
|
||||||
.await
|
.await
|
||||||
.unwrap();
|
.unwrap();
|
||||||
repo.insert(trec_full("c2", "todo", Some("parent"), "todo"))
|
repo.insert(trec_full("c2", "todo", Some("parent"), TaskStatus::Todo))
|
||||||
.await
|
.await
|
||||||
.unwrap();
|
.unwrap();
|
||||||
repo.insert(trec_full("c3", "todo", Some("parent"), "in_progress"))
|
repo.insert(trec_full("c3", "todo", Some("parent"), TaskStatus::InProgress))
|
||||||
.await
|
.await
|
||||||
.unwrap();
|
.unwrap();
|
||||||
repo.insert(trec_full("c4", "todo", Some("parent"), "done"))
|
repo.insert(trec_full("c4", "todo", Some("parent"), TaskStatus::Done))
|
||||||
.await
|
.await
|
||||||
.unwrap();
|
.unwrap();
|
||||||
|
|
||||||
@@ -841,7 +841,7 @@ mod tests {
|
|||||||
repo.insert(trec_full("parent", "todo", None, TaskStatus::Todo))
|
repo.insert(trec_full("parent", "todo", None, TaskStatus::Todo))
|
||||||
.await
|
.await
|
||||||
.unwrap();
|
.unwrap();
|
||||||
let ok = repo.set_status_for_aggregation("parent", "in_progress");
|
let ok = repo.set_status_for_aggregation("parent", "in_progress").await.unwrap();
|
||||||
assert!(ok, "应命中写入");
|
assert!(ok, "应命中写入");
|
||||||
let after = repo.get_by_id("parent").await.unwrap().unwrap();
|
let after = repo.get_by_id("parent").await.unwrap().unwrap();
|
||||||
assert_eq!(after.status.as_str(), "in_progress", "status 应被聚合写入更新");
|
assert_eq!(after.status.as_str(), "in_progress", "status 应被聚合写入更新");
|
||||||
|
|||||||
@@ -45,7 +45,7 @@ pub fn run(conn: &Connection) -> Result<()> {
|
|||||||
// 什么数据库、Redis 在哪、有没有 MQ"的基础设施上下文。
|
// 什么数据库、Redis 在哪、有没有 MQ"的基础设施上下文。
|
||||||
// V33 = 审批重启恢复:ai_conversations 加 pending_approvals TEXT 列,持久化挂起审批快照,
|
// V33 = 审批重启恢复:ai_conversations 加 pending_approvals TEXT 列,持久化挂起审批快照,
|
||||||
// 重启后从 DB 恢复 pending_approvals 内存态,使待审批不丢。
|
// 重启后从 DB 恢复 pending_approvals 内存态,使待审批不丢。
|
||||||
let steps: [(i32, fn(&Connection) -> Result<()>); 35] = [
|
let steps: [(i32, fn(&Connection) -> Result<()>); 36] = [
|
||||||
(1, migrate_v1),
|
(1, migrate_v1),
|
||||||
(2, migrate_v2),
|
(2, migrate_v2),
|
||||||
(3, migrate_v3),
|
(3, migrate_v3),
|
||||||
@@ -81,6 +81,7 @@ pub fn run(conn: &Connection) -> Result<()> {
|
|||||||
(33, migrate_v33),
|
(33, migrate_v33),
|
||||||
(34, migrate_v34),
|
(34, migrate_v34),
|
||||||
(35, migrate_v35),
|
(35, migrate_v35),
|
||||||
|
(36, migrate_v36),
|
||||||
];
|
];
|
||||||
|
|
||||||
for (version, migrate_fn) in steps {
|
for (version, migrate_fn) in steps {
|
||||||
@@ -1030,6 +1031,90 @@ fn migrate_v35(conn: &Connection) -> Result<()> {
|
|||||||
Ok(())
|
Ok(())
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/// V36:多 Agent 并行执行数据层(ai_plans/ai_subtasks/ai_conflicts 3 新表 +
|
||||||
|
/// ai_messages/ai_tool_executions 加 subtask_id 列)。
|
||||||
|
/// 设计依据:docs/02-架构设计/专项设计/多Agent并行执行与仲裁合并设计-2026-07-01.md §二
|
||||||
|
fn migrate_v36(conn: &Connection) -> Result<()> {
|
||||||
|
// 1. ai_messages 加 subtask_id 列(消息归属子任务,NULL=单 Agent 时期)
|
||||||
|
if !column_exists(conn, "ai_messages", "subtask_id") {
|
||||||
|
conn.execute("ALTER TABLE ai_messages ADD COLUMN subtask_id TEXT", [])?;
|
||||||
|
tracing::info!("v36: ai_messages 加 subtask_id 列");
|
||||||
|
}
|
||||||
|
|
||||||
|
// 2. ai_tool_executions 加 subtask_id 列(工具调用归属子任务)
|
||||||
|
if !column_exists(conn, "ai_tool_executions", "subtask_id") {
|
||||||
|
conn.execute("ALTER TABLE ai_tool_executions ADD COLUMN subtask_id TEXT", [])?;
|
||||||
|
tracing::info!("v36: ai_tool_executions 加 subtask_id 列");
|
||||||
|
}
|
||||||
|
|
||||||
|
// 3. ai_plans 表(Plan 生命周期:用户消息触发→拆解→执行→合并→完成)
|
||||||
|
conn.execute(
|
||||||
|
"CREATE TABLE IF NOT EXISTS ai_plans (
|
||||||
|
id TEXT PRIMARY KEY,
|
||||||
|
conversation_id TEXT NOT NULL,
|
||||||
|
user_message_id TEXT,
|
||||||
|
status TEXT NOT NULL DEFAULT 'planning',
|
||||||
|
subtask_count INTEGER NOT NULL DEFAULT 0,
|
||||||
|
created_at TEXT NOT NULL,
|
||||||
|
completed_at TEXT
|
||||||
|
)",
|
||||||
|
[],
|
||||||
|
)?;
|
||||||
|
conn.execute(
|
||||||
|
"CREATE INDEX IF NOT EXISTS idx_ai_plans_conv ON ai_plans(conversation_id, created_at)",
|
||||||
|
[],
|
||||||
|
)?;
|
||||||
|
|
||||||
|
// 4. ai_subtasks 表(SubTask 状态 + DAG 层级 + Git worktree 分支)
|
||||||
|
conn.execute(
|
||||||
|
"CREATE TABLE IF NOT EXISTS ai_subtasks (
|
||||||
|
id TEXT PRIMARY KEY,
|
||||||
|
plan_id TEXT NOT NULL REFERENCES ai_plans(id),
|
||||||
|
persona_id TEXT,
|
||||||
|
intent TEXT NOT NULL DEFAULT '',
|
||||||
|
status TEXT NOT NULL DEFAULT 'pending',
|
||||||
|
layer INTEGER NOT NULL DEFAULT 0,
|
||||||
|
deps TEXT,
|
||||||
|
branch TEXT,
|
||||||
|
created_at TEXT NOT NULL,
|
||||||
|
completed_at TEXT
|
||||||
|
)",
|
||||||
|
[],
|
||||||
|
)?;
|
||||||
|
conn.execute(
|
||||||
|
"CREATE INDEX IF NOT EXISTS idx_ai_subtasks_plan ON ai_subtasks(plan_id, layer)",
|
||||||
|
[],
|
||||||
|
)?;
|
||||||
|
|
||||||
|
// 5. ai_conflicts 表(合并冲突:同文件路径 + 跨文件语义冲突)
|
||||||
|
conn.execute(
|
||||||
|
"CREATE TABLE IF NOT EXISTS ai_conflicts (
|
||||||
|
id TEXT PRIMARY KEY,
|
||||||
|
plan_id TEXT NOT NULL REFERENCES ai_plans(id),
|
||||||
|
file_path TEXT NOT NULL DEFAULT '',
|
||||||
|
conflict_type TEXT NOT NULL DEFAULT 'file',
|
||||||
|
subtask_a TEXT,
|
||||||
|
subtask_b TEXT,
|
||||||
|
diff_a TEXT,
|
||||||
|
diff_b TEXT,
|
||||||
|
resolution TEXT NOT NULL DEFAULT 'pending',
|
||||||
|
resolved_by TEXT,
|
||||||
|
created_at TEXT NOT NULL,
|
||||||
|
resolved_at TEXT
|
||||||
|
)",
|
||||||
|
[],
|
||||||
|
)?;
|
||||||
|
conn.execute(
|
||||||
|
"CREATE INDEX IF NOT EXISTS idx_ai_conflicts_plan ON ai_conflicts(plan_id, resolution)",
|
||||||
|
[],
|
||||||
|
)?;
|
||||||
|
|
||||||
|
tracing::info!("v36: 建 ai_plans/ai_subtasks/ai_conflicts 表 + subtask_id 列(多 Agent 并行执行数据层)");
|
||||||
|
conn.execute("INSERT INTO schema_version (version) VALUES (?)", [36])?;
|
||||||
|
tracing::info!("迁移 v36 完成");
|
||||||
|
Ok(())
|
||||||
|
}
|
||||||
|
|
||||||
/// V21 建表 SQL — 消息拆分存储 ai_messages 表
|
/// V21 建表 SQL — 消息拆分存储 ai_messages 表
|
||||||
///
|
///
|
||||||
/// 与 V9_SQL 中的 ai_messages 镜像(V9 给新库,此 const 给老库 V21 迁移用 IF NOT EXISTS)。
|
/// 与 V9_SQL 中的 ai_messages 镜像(V9 给新库,此 const 给老库 V21 迁移用 IF NOT EXISTS)。
|
||||||
|
|||||||
@@ -497,3 +497,62 @@ pub struct ModuleDependencyRecord {
|
|||||||
pub label: Option<String>,
|
pub label: Option<String>,
|
||||||
pub created_at: String,
|
pub created_at: String,
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// ============================================================
|
||||||
|
// 多 Agent 并行执行模型 (V36)
|
||||||
|
// ============================================================
|
||||||
|
|
||||||
|
/// Plan 记录(一次用户消息触发的并行执行计划)
|
||||||
|
#[derive(Debug, Clone, Serialize, Deserialize)]
|
||||||
|
pub struct PlanRecord {
|
||||||
|
pub id: String,
|
||||||
|
pub conversation_id: String,
|
||||||
|
/// 触发 Plan 的用户消息 id
|
||||||
|
pub user_message_id: Option<String>,
|
||||||
|
/// planning / executing / merging / done / error
|
||||||
|
pub status: String,
|
||||||
|
pub subtask_count: i64,
|
||||||
|
pub created_at: String,
|
||||||
|
pub completed_at: Option<String>,
|
||||||
|
}
|
||||||
|
|
||||||
|
/// SubTask 记录(Plan 拆解出的子任务)
|
||||||
|
#[derive(Debug, Clone, Serialize, Deserialize)]
|
||||||
|
pub struct SubTaskRecord {
|
||||||
|
pub id: String,
|
||||||
|
pub plan_id: String,
|
||||||
|
/// 分配的人设 id(coder/reviewer/architect/tester/analyst)
|
||||||
|
pub persona_id: Option<String>,
|
||||||
|
/// 子任务意图描述
|
||||||
|
pub intent: String,
|
||||||
|
/// pending / running / done / error
|
||||||
|
pub status: String,
|
||||||
|
/// DAG 层级(0 起)
|
||||||
|
pub layer: i64,
|
||||||
|
/// 依赖的 SubTask id 列表(JSON 数组字符串)
|
||||||
|
pub deps: Option<String>,
|
||||||
|
/// Git worktree 分支名(NULL=非 Git 工程降级串行)
|
||||||
|
pub branch: Option<String>,
|
||||||
|
pub created_at: String,
|
||||||
|
pub completed_at: Option<String>,
|
||||||
|
}
|
||||||
|
|
||||||
|
/// 冲突记录(合并时检测到的文件/语义冲突)
|
||||||
|
#[derive(Debug, Clone, Serialize, Deserialize)]
|
||||||
|
pub struct ConflictRecord {
|
||||||
|
pub id: String,
|
||||||
|
pub plan_id: String,
|
||||||
|
pub file_path: String,
|
||||||
|
/// file(同文件路径) / semantic(编译失败)
|
||||||
|
pub conflict_type: String,
|
||||||
|
pub subtask_a: Option<String>,
|
||||||
|
pub subtask_b: Option<String>,
|
||||||
|
pub diff_a: Option<String>,
|
||||||
|
pub diff_b: Option<String>,
|
||||||
|
/// pending / a / b / merged / manual
|
||||||
|
pub resolution: String,
|
||||||
|
/// reviewer / user
|
||||||
|
pub resolved_by: Option<String>,
|
||||||
|
pub created_at: String,
|
||||||
|
pub resolved_at: Option<String>,
|
||||||
|
}
|
||||||
|
|||||||
Reference in New Issue
Block a user