From ae6d3d0043336984d8acaf6c656e10d89cf44e09 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=E7=BB=9D=E5=B0=98?= <237809796@qq.com> Date: Wed, 1 Jul 2026 21:50:16 +0800 Subject: [PATCH] =?UTF-8?q?=E6=96=B0=E5=A2=9E:=20=E5=A4=9AAgent=E6=95=B0?= =?UTF-8?q?=E6=8D=AE=E5=B1=82(V36=E8=BF=81=E7=A7=BB+3=E6=96=B0=E8=A1=A8+Re?= =?UTF-8?q?po=20CRUD+7=E4=B8=AA=E6=B5=8B=E8=AF=95=E5=85=A8=E7=BB=BF)?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - 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) --- Batch.md | 2 + crates/df-storage/src/crud/mod.rs | 5 +- crates/df-storage/src/crud/plan_repo.rs | 579 ++++++++++++++++++++++++ crates/df-storage/src/crud/task_repo.rs | 10 +- crates/df-storage/src/migrations.rs | 87 +++- crates/df-storage/src/models.rs | 59 +++ 6 files changed, 735 insertions(+), 7 deletions(-) create mode 100644 crates/df-storage/src/crud/plan_repo.rs diff --git a/Batch.md b/Batch.md index 7a9cdaa..32598d0 100644 --- a/Batch.md +++ b/Batch.md @@ -437,6 +437,8 @@ ### 数据层 + 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 | diff --git a/crates/df-storage/src/crud/mod.rs b/crates/df-storage/src/crud/mod.rs index 6ce80b5..2459a24 100644 --- a/crates/df-storage/src/crud/mod.rs +++ b/crates/df-storage/src/crud/mod.rs @@ -17,6 +17,9 @@ //! re-export(`pub use ...::*`)保持 `df_storage::crud::XxxRepo` / //! `df_storage::crud::is_allowed_column` 路径不变,**调用方零改动**。 +mod module_dependency_repo; +mod plan_repo; + mod conversation_repo; mod idea_eval_repo; mod idea_repo; @@ -28,7 +31,6 @@ mod project_service_repo; mod settings; mod task_link_repo; mod task_repo; -mod module_dependency_repo; pub use conversation_repo::*; pub use idea_eval_repo::*; @@ -42,6 +44,7 @@ pub use settings::*; pub use task_link_repo::*; pub use task_repo::*; pub use module_dependency_repo::*; +pub use plan_repo::*; // ============================================================ // 辅助宏 — 消除多个 Repo 的重复样板 diff --git a/crates/df-storage/src/crud/plan_repo.rs b/crates/df-storage/src/crud/plan_repo.rs new file mode 100644 index 0000000..51c0233 --- /dev/null +++ b/crates/df-storage/src/crud/plan_repo.rs @@ -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 = std::result::Result; + +fn plan_from_row(row: &Row<'_>) -> rusqlite::Result { + 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 { + 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 { + 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>, +} + +impl PlanRepo { + pub fn new(db: &Database) -> Self { + Self { conn: db.conn() } + } + + pub async fn insert(&self, record: PlanRecord) -> Result { + 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> { + 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 { + 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> { + 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>, +} + +impl SubTaskRepo { + pub fn new(db: &Database) -> Self { + Self { conn: db.conn() } + } + + pub async fn insert(&self, record: SubTaskRecord) -> Result { + 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> { + 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 { + 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>, +} + +impl ConflictRepo { + pub fn new(db: &Database) -> Self { + Self { conn: db.conn() } + } + + pub async fn insert(&self, record: ConflictRecord) -> Result { + 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> { + 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 { + 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"); + } +} diff --git a/crates/df-storage/src/crud/task_repo.rs b/crates/df-storage/src/crud/task_repo.rs index 0be7cd1..03cdca0 100644 --- a/crates/df-storage/src/crud/task_repo.rs +++ b/crates/df-storage/src/crud/task_repo.rs @@ -784,16 +784,16 @@ mod tests { 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")) + repo.insert(trec_full("c1", "todo", Some("parent"), TaskStatus::Todo)) .await .unwrap(); - repo.insert(trec_full("c2", "todo", Some("parent"), "todo")) + repo.insert(trec_full("c2", "todo", Some("parent"), TaskStatus::Todo)) .await .unwrap(); - repo.insert(trec_full("c3", "todo", Some("parent"), "in_progress")) + repo.insert(trec_full("c3", "todo", Some("parent"), TaskStatus::InProgress)) .await .unwrap(); - repo.insert(trec_full("c4", "todo", Some("parent"), "done")) + repo.insert(trec_full("c4", "todo", Some("parent"), TaskStatus::Done)) .await .unwrap(); @@ -841,7 +841,7 @@ mod tests { repo.insert(trec_full("parent", "todo", None, TaskStatus::Todo)) .await .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, "应命中写入"); let after = repo.get_by_id("parent").await.unwrap().unwrap(); assert_eq!(after.status.as_str(), "in_progress", "status 应被聚合写入更新"); diff --git a/crates/df-storage/src/migrations.rs b/crates/df-storage/src/migrations.rs index 430cfe6..0e51739 100644 --- a/crates/df-storage/src/migrations.rs +++ b/crates/df-storage/src/migrations.rs @@ -45,7 +45,7 @@ pub fn run(conn: &Connection) -> Result<()> { // 什么数据库、Redis 在哪、有没有 MQ"的基础设施上下文。 // V33 = 审批重启恢复:ai_conversations 加 pending_approvals TEXT 列,持久化挂起审批快照, // 重启后从 DB 恢复 pending_approvals 内存态,使待审批不丢。 - let steps: [(i32, fn(&Connection) -> Result<()>); 35] = [ + let steps: [(i32, fn(&Connection) -> Result<()>); 36] = [ (1, migrate_v1), (2, migrate_v2), (3, migrate_v3), @@ -81,6 +81,7 @@ pub fn run(conn: &Connection) -> Result<()> { (33, migrate_v33), (34, migrate_v34), (35, migrate_v35), + (36, migrate_v36), ]; for (version, migrate_fn) in steps { @@ -1030,6 +1031,90 @@ fn migrate_v35(conn: &Connection) -> Result<()> { 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 表 /// /// 与 V9_SQL 中的 ai_messages 镜像(V9 给新库,此 const 给老库 V21 迁移用 IF NOT EXISTS)。 diff --git a/crates/df-storage/src/models.rs b/crates/df-storage/src/models.rs index 8c95554..3cbf453 100644 --- a/crates/df-storage/src/models.rs +++ b/crates/df-storage/src/models.rs @@ -497,3 +497,62 @@ pub struct ModuleDependencyRecord { pub label: Option, 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, + /// planning / executing / merging / done / error + pub status: String, + pub subtask_count: i64, + pub created_at: String, + pub completed_at: Option, +} + +/// 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, + /// 子任务意图描述 + pub intent: String, + /// pending / running / done / error + pub status: String, + /// DAG 层级(0 起) + pub layer: i64, + /// 依赖的 SubTask id 列表(JSON 数组字符串) + pub deps: Option, + /// Git worktree 分支名(NULL=非 Git 工程降级串行) + pub branch: Option, + pub created_at: String, + pub completed_at: Option, +} + +/// 冲突记录(合并时检测到的文件/语义冲突) +#[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, + pub subtask_b: Option, + pub diff_a: Option, + pub diff_b: Option, + /// pending / a / b / merged / manual + pub resolution: String, + /// reviewer / user + pub resolved_by: Option, + pub created_at: String, + pub resolved_at: Option, +}