From 9eb2995a74e90f2669362d816ef9f0f01c6d5aa4 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=E7=BB=9D=E5=B0=98?= <237809796@qq.com> Date: Sat, 8 Aug 2026 21:24:53 +0800 Subject: [PATCH] =?UTF-8?q?=E6=96=B0=E5=A2=9E:=20=E5=B7=A5=E5=85=B7?= =?UTF-8?q?=E5=B7=A5=E4=BD=9C=E6=B5=81=E8=A1=A5=E5=85=A8(diff=5Ffiles?= =?UTF-8?q?=E5=B7=A5=E5=85=B7=20+=20=E6=8A=80=E8=83=BD=E6=B8=85=E5=8D=95?= =?UTF-8?q?=E6=B3=A8=E5=85=A5=20+=20=E5=B7=A5=E4=BD=9C=E6=B5=81=E8=BF=9B?= =?UTF-8?q?=E5=BA=A6=E5=85=B1=E4=BA=AB=20+=20human=E7=AB=AF=E5=88=B0?= =?UTF-8?q?=E7=AB=AF=E6=B5=8B=E8=AF=95=20+=20=E7=9F=A5=E8=AF=86=E5=BA=93MC?= =?UTF-8?q?P=E5=B7=A5=E5=85=B7)?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- crates/df-ai/src/namespace_store.rs | 1 + crates/df-mcp/src/tools.rs | 254 +++++++++++++++++- crates/df-workflow/tests/human_workflow.rs | 235 ++++++++++++++++ src-tauri/src/commands/ai/agentic/mod.rs | 56 ++++ src-tauri/src/commands/ai/tool_registry.rs | 15 +- src-tauri/src/commands/ai/tools/file.rs | 72 +++++ src/api/types.ts | 3 + .../workflow/WorkflowDagDisplay.vue | 3 +- src/i18n/en/projectDetail.ts | 2 + src/i18n/zh-CN/projectDetail.ts | 2 + src/stores/project.ts | 2 + src/stores/project/workflow.ts | 81 ++++++ src/views/ProjectDetail.vue | 63 ++++- src/views/TaskDetail.vue | 111 +++----- 14 files changed, 807 insertions(+), 93 deletions(-) create mode 100644 crates/df-workflow/tests/human_workflow.rs diff --git a/crates/df-ai/src/namespace_store.rs b/crates/df-ai/src/namespace_store.rs index fa21051..f314c04 100644 --- a/crates/df-ai/src/namespace_store.rs +++ b/crates/df-ai/src/namespace_store.rs @@ -196,6 +196,7 @@ mod tests { fn should_use_namespace_always_large_tool() { assert!(should_use_namespace("short", "read_file")); assert!(should_use_namespace("short", "list_directory")); + assert!(should_use_namespace("short", "diff_files")); } #[test] diff --git a/crates/df-mcp/src/tools.rs b/crates/df-mcp/src/tools.rs index addc70e..a91703e 100644 --- a/crates/df-mcp/src/tools.rs +++ b/crates/df-mcp/src/tools.rs @@ -15,9 +15,9 @@ // name→id 解析:src-tauri 有机制层解析(audit/mod.rs auto_resolve),MCP 面暂不同步。 use std::sync::{Arc, OnceLock}; -use df_storage::crud::{IdeaQuery, IdeaRepo, ProjectQuery, ProjectRepo, TaskQuery, TaskRepo}; +use df_storage::crud::{IdeaQuery, IdeaRepo, KnowledgeRepo, ProjectQuery, ProjectRepo, TaskQuery, TaskRepo}; use df_storage::db::Database; -use df_storage::models::{IdeaRecord, ProjectRecord, TaskRecord}; +use df_storage::models::{IdeaRecord, KnowledgeRecord, ProjectRecord, TaskRecord}; use df_types::types::{IdeaStatus, ProjectStatus, TaskStatus, new_id}; use futures::future::BoxFuture; use serde_json::{json, Value}; @@ -123,6 +123,22 @@ pub fn all_tools() -> &'static Vec<&'static ToolSpec> { spec("score_idea", "评分并写库(Medium 风险+审计日志):对想法做启发式评估,把 scores 写回 DB 并返回更新后的记录", object_schema(json!({"id": str_field("想法 ID"), "expected_updated_at": int_field("乐观锁版本(可空):上次读取到的 updated_at 毫秒时间戳,不一致则拒绝写入")}), &["id"]), Medium, score_idea), // ─── 工作流(High) ─── spec("run_workflow", "触发工作流——High 风险,默认拒绝,请在 DevFlow 应用内执行", object_schema(json!({"project_id": str_field("项目 ID"), "task_id": opt_str_field("任务 ID(可空)")}), &["project_id"]), High, run_workflow), + // ─── 知识库 ─── + spec("search_knowledge", "检索知识库:按关键词 LIKE 匹配 title/content(可选 kind/limit);传 query_embedding 数组则走向量检索(余弦相似度 top-N)", object_schema(json!({ + "query": str_field("检索关键词"), + "kind": opt_str_field("知识类型过滤(可空):review_rule/prompt_template/pitfall/architecture_pattern/diagnosis/deployment_note/workflow_optimization"), + "limit": int_field("返回上限(可空,默认 5,上限 10)"), + "query_embedding": json!({ "type": "array", "items": { "type": "number" }, "description": "查询向量(可空,传则走向量检索)" }) + }), &["query"]), Low, search_knowledge), + spec("insert_knowledge", "新增知识条目(Medium 风险,默认允许+审计日志;状态恒为 candidate,审核发布后才参与检索)", object_schema(json!({ + "kind": str_field("知识类型:review_rule/prompt_template/pitfall/architecture_pattern/diagnosis/deployment_note/workflow_optimization"), + "title": str_field("标题"), + "content": str_field("内容"), + "tags": opt_str_field("标签 JSON 数组字符串(可空)"), + "source_project": opt_str_field("来源项目(可空)"), + "source_ref": opt_str_field("来源引用(可空,如 conv:{id})"), + "confidence": opt_str_field("置信度(可空:high/medium/low)") + }), &["kind", "title", "content"]), Medium, insert_knowledge), // ─── 回收站 ─── spec("list_trash", "列出回收站(deleted_at IS NOT NULL 的项目与任务)", object_schema(json!({}), &[]), Low, list_trash), spec("restore_project", "从回收站恢复项目(Medium 风险+审计日志)", object_schema(json!({"id": str_field("项目 ID")}), &["id"]), Medium, restore_project), @@ -980,6 +996,126 @@ fn restore_project(ctx: &Ctx, args: Value) -> BoxFuture<'static, CallToolResult> }) } +// ============================================================ +// handler 实现 — 知识库(复用 df-storage KnowledgeRepo) +// ============================================================ +// +// 分层存储(Tier 2/3:热/温/冷知识分级)需单独设计,本批只接 Tier 1 扁平 +// KnowledgeRepo(单表 + 向量列),不引入存储分层。 + +/// 7 种合法知识类型(对齐 migrations.rs V7 建表注释)。 +const KNOWLEDGE_KINDS: &[&str] = &[ + "review_rule", "prompt_template", "pitfall", "architecture_pattern", + "diagnosis", "deployment_note", "workflow_optimization", +]; + +/// 检索知识(Low 只读):默认关键词 LIKE(title/content),可选 kind/limit; +/// 传 query_embedding(数组)则走向量检索(余弦 top-N),返回带相似度。 +/// 向量由调用方生成(MCP 无 AI provider 上下文),与 GUI 的 hybrid_search 共用 +/// KnowledgeRepo::search_vector,结果一致只做列裁剪。 +fn search_knowledge(ctx: &Ctx, args: Value) -> BoxFuture<'static, CallToolResult> { + let db = ctx.db.clone(); + let query = match arg_str(&args, "query") { + Ok(v) => v, + Err(r) => return Box::pin(std::future::ready(r)), + }; + let kind = args.get("kind").and_then(|v| v.as_str()).map(|s| s.to_owned()); + let limit = args + .get("limit") + .and_then(|v| v.as_i64()) + .map(|i| i.max(1) as usize) + .unwrap_or(5) + .min(10); + let query_vec: Option> = args + .get("query_embedding") + .and_then(|v| v.as_array()) + .map(|arr| arr.iter().filter_map(|n| n.as_f64()).map(|f| f as f32).collect()); + Box::pin(async move { + let repo = KnowledgeRepo::new(&db); + match query_vec { + Some(vec) => { + if vec.is_empty() { + return CallToolResult::error("query_embedding 不能为空数组"); + } + match repo.search_vector(&vec, limit).await { + Ok(hits) => { + let hits: Vec = hits + .into_iter() + .map(|(rec, score)| json!({ + "score": (score * 1000.0).round() / 1000.0, + "knowledge": rec + })) + .collect(); + json_ok(json!({ "count": hits.len(), "hits": hits })) + } + Err(e) => err_str(e), + } + } + None => match repo.search(&query, kind.as_deref(), limit).await { + Ok(list) => json_ok(json!({ "count": list.len(), "knowledge": list })), + Err(e) => err_str(e), + }, + } + }) +} + +/// 新增知识条目(Medium 写库):状态恒为 candidate(审核发布后才参与检索, +/// 对齐 GUI knowledge_create 语义),不自动生成嵌入(嵌入由发布链路处理)。 +fn insert_knowledge(ctx: &Ctx, args: Value) -> BoxFuture<'static, CallToolResult> { + let db = ctx.db.clone(); + let kind = match arg_str(&args, "kind") { + Ok(v) => v, + Err(r) => return Box::pin(std::future::ready(r)), + }; + let title = match arg_str(&args, "title") { + Ok(v) => v, + Err(r) => return Box::pin(std::future::ready(r)), + }; + let content = match arg_str(&args, "content") { + Ok(v) => v, + Err(r) => return Box::pin(std::future::ready(r)), + }; + let tags = args.get("tags").and_then(|v| v.as_str()).map(|s| s.to_owned()); + let source_project = args.get("source_project").and_then(|v| v.as_str()).map(|s| s.to_owned()); + let source_ref = args.get("source_ref").and_then(|v| v.as_str()).map(|s| s.to_owned()); + let confidence = args.get("confidence").and_then(|v| v.as_str()).map(|s| s.to_owned()); + medium_audit("insert_knowledge", &title); + Box::pin(async move { + if !KNOWLEDGE_KINDS.contains(&kind.as_str()) { + return CallToolResult::error(format!( + "非法知识类型: {kind}, 有效值: {}", + KNOWLEDGE_KINDS.join("/") + )); + } + if title.trim().is_empty() || content.trim().is_empty() { + return CallToolResult::error("title/content 不能为空"); + } + let now = now_millis(); + let rec = KnowledgeRecord { + id: new_id(), + kind, + title, + content, + tags, + status: "candidate".to_string(), + confidence, + reuse_count: 0, + verified: false, + source_project, + source_ref, + reasoning: None, + embedding_status: None, + created_at: now.clone(), + updated_at: now, + }; + let repo = KnowledgeRepo::new(&db); + match repo.insert(rec.clone()).await { + Ok(id) => json_ok(json!({ "id": id, "knowledge": rec })), + Err(e) => err_str(e), + } + }) +} + // ============================================================ // 单测:evaluate_idea(只读,不写库)/ score_idea(写库)/ 风险契约 // ============================================================ @@ -1079,6 +1215,30 @@ mod tests { repo.insert(rec).await.unwrap() } + /// 插入一条知识,返回 id(search_knowledge 测试用;status=published 才进检索) + async fn seed_knowledge(ctx: &Ctx, title: &str, content: &str, kind: &str) -> String { + let repo = KnowledgeRepo::new(&ctx.db); + let now = now_millis(); + let rec = KnowledgeRecord { + id: new_id(), + kind: kind.to_owned(), + title: title.to_owned(), + content: content.to_owned(), + tags: None, + status: "published".to_string(), + confidence: Some("high".to_owned()), + reuse_count: 0, + verified: true, + source_project: None, + source_ref: None, + reasoning: None, + embedding_status: None, + created_at: now.clone(), + updated_at: now, + }; + repo.insert(rec).await.unwrap() + } + /// 读当前 DB 中的 idea.scores(原始字符串) async fn db_scores(ctx: &Ctx, id: &str) -> Option { IdeaRepo::new(&ctx.db) @@ -1557,4 +1717,94 @@ mod tests { let stack = v["project"]["stack"].as_str().unwrap_or_default(); assert!(stack.contains("rust"), "应探测到 rust 技术栈,实际 stack: {stack}"); } + + // ── 知识库工具(search_knowledge Low / insert_knowledge Medium)──────────── + + /// search_knowledge:关键词命中已发布知识,返回记录列表(不写库,只读契约)。 + #[tokio::test] + async fn search_knowledge_keyword_returns_hits() { + let ctx = test_ctx().await; + seed_knowledge(&ctx, "Rust 异步模型", "tokio 运行时与并发", "pitfall").await; + + let r = search_knowledge(&ctx, json!({ "query": "tokio" })).await; + assert!(r.is_error.is_none(), "{:?}", text_of(&r)); + let v = json_of(&r); + assert_eq!(v["count"], 1); + assert_eq!(v["knowledge"][0]["title"], "Rust 异步模型"); + } + + /// search_knowledge:缺 query 必填参数 → 报错。 + #[tokio::test] + async fn search_knowledge_missing_query_errors() { + let ctx = test_ctx().await; + let r = search_knowledge(&ctx, json!({})).await; + assert_eq!(r.is_error, Some(true)); + assert!(text_of(&r).contains("缺少必填参数")); + } + + /// search_knowledge:传入 query_embedding → 向量检索,返回带 score 的命中。 + #[tokio::test] + async fn search_knowledge_vector_with_embedding() { + let ctx = test_ctx().await; + let id = seed_knowledge(&ctx, "Rust 异步", "tokio 并发模型", "pitfall").await; + // 写一条向量(f32 数组,维度 3),search_vector 才能命中 + KnowledgeRepo::new(&ctx.db).set_embedding(&id, &[0.1, 0.2, 0.3]).await.unwrap(); + + let r = search_knowledge(&ctx, json!({ + "query": "ignored", // 有 embedding 时 query 仅作占位,检索走向量 + "query_embedding": [0.1, 0.2, 0.3] + })).await; + assert!(r.is_error.is_none(), "{:?}", text_of(&r)); + let v = json_of(&r); + assert_eq!(v["count"], 1); + assert!(v["hits"][0]["score"].is_number()); + assert_eq!(v["hits"][0]["knowledge"]["id"], id); + } + + /// insert_knowledge:创建成功,状态恒 candidate(待审核),DB 可读回。 + #[tokio::test] + async fn insert_knowledge_creates_candidate() { + let ctx = test_ctx().await; + let r = insert_knowledge(&ctx, json!({ + "kind": "pitfall", + "title": "MCP 超时", + "content": "审批响应须带 execution_id 匹配", + "tags": "[\"mcp\",\"workflow\"]" + })).await; + assert!(r.is_error.is_none(), "{:?}", text_of(&r)); + let v = json_of(&r); + let id = v["id"].as_str().unwrap().to_string(); + assert_eq!(v["knowledge"]["status"], "candidate"); + assert_eq!(v["knowledge"]["reuse_count"], 0); + let persisted = KnowledgeRepo::new(&ctx.db).get_by_id(&id).await.unwrap(); + assert!(persisted.is_some(), "DB 应能读回新建知识"); + assert_eq!(persisted.unwrap().status, "candidate"); + } + + /// insert_knowledge:非法 kind / 空 title → 报错(不落库)。 + #[tokio::test] + async fn insert_knowledge_invalid_input_errors() { + let ctx = test_ctx().await; + let bad_kind = insert_knowledge(&ctx, json!({ "kind": "nope", "title": "t", "content": "c" })).await; + assert_eq!(bad_kind.is_error, Some(true)); + assert!(text_of(&bad_kind).contains("非法知识类型")); + + let empty_title = insert_knowledge(&ctx, json!({ "kind": "pitfall", "title": " ", "content": "c" })).await; + assert_eq!(empty_title.is_error, Some(true)); + assert!(text_of(&empty_title).contains("不能为空")); + } + + /// 风险契约:search_knowledge=Low(只读,read-only 放行),insert_knowledge=Medium(写库,read-only 拒)。 + #[test] + fn knowledge_tools_risk_contract() { + let search = find("search_knowledge").expect("search_knowledge 必须注册"); + let insert = find("insert_knowledge").expect("insert_knowledge 必须注册"); + assert_eq!(search.risk, RiskLevel::Low, "search_knowledge 必须 Low(只读契约)"); + assert_eq!(insert.risk, RiskLevel::Medium, "insert_knowledge 必须 Medium(写库 → read-only 拒)"); + + assert!(visible_for_test(true, "search_knowledge")); + assert!(!visible_for_test(true, "insert_knowledge")); + assert!(visible_for_test(false, "search_knowledge")); + assert!(visible_for_test(false, "insert_knowledge")); + } } diff --git a/crates/df-workflow/tests/human_workflow.rs b/crates/df-workflow/tests/human_workflow.rs new file mode 100644 index 0000000..50e075f --- /dev/null +++ b/crates/df-workflow/tests/human_workflow.rs @@ -0,0 +1,235 @@ +//! human 节点端到端集成测试 — DagDef → NodeRegistry.build_dag → DagExecutor.run +//! +//! df-workflow 不依赖 df-nodes(反向依赖),此处用等价阻塞节点模拟 HumanNode 的 +//! 「发审批请求 → 挂起等待 → 外部 approve → 返回结果」链路,覆盖 +//! DagDef 序列化 → 注册表 → 执行器全链路 + 事件序列断言(R6 send 缺 await / +//! R7 契约失配的存活土壤)。 + +use std::time::Duration; + +use async_trait::async_trait; +use serde_json::json; + +use df_types::events::{SelectType, WorkflowEvent}; +use df_types::types::NodeStatus; +use df_workflow::dag_def::DagDef; +use df_workflow::eventbus::{EventBus, EventSubscriber}; +use df_workflow::executor::DagExecutor; +use df_workflow::node::{Node, NodeContext, NodeOutput, NodeResult, NodeSchema}; +use df_workflow::registry::NodeRegistry; + +/// 前驱节点:sleep 后返回空输出(仿 executor 单测 SleepNode)。 +struct SleepNode { + sleep_ms: u64, +} + +#[async_trait] +impl Node for SleepNode { + async fn execute(&self, _ctx: NodeContext) -> NodeResult { + tokio::time::sleep(Duration::from_millis(self.sleep_ms)).await; + Ok(NodeOutput::empty()) + } + + fn schema(&self) -> NodeSchema { + NodeSchema { + params: json!(null), + output: json!(null), + } + } + + fn node_type(&self) -> &str { + "sleep" + } +} + +/// 模拟 HumanNode 的阻塞审批节点:先 subscribe 再发 HumanApprovalRequest(broadcast 不回放), +/// select! 等待匹配 execution_id + node_id 的 HumanApprovalResponse,收到则返回 decision。 +struct ApprovalNode; + +#[async_trait] +impl Node for ApprovalNode { + async fn execute(&self, ctx: NodeContext) -> NodeResult { + let mut rx = ctx.event_bus.subscribe(); + let _ = ctx + .event_bus + .send(WorkflowEvent::HumanApprovalRequest { + execution_id: ctx.execution_id.clone(), + node_id: ctx.node_id.clone(), + title: "确认发布".to_string(), + description: String::new(), + options: vec!["同意".into(), "拒绝".into()], + select_type: SelectType::Single, + }) + .await; + let timeout_secs = ctx + .config + .get("timeout_secs") + .and_then(|v| v.as_u64()) + .unwrap_or(2); + let deadline = tokio::time::Instant::now() + Duration::from_secs(timeout_secs); + loop { + tokio::select! { + recv = rx.recv() => match recv { + Ok(WorkflowEvent::HumanApprovalResponse { + execution_id, node_id, decision, decisions, .. + }) if execution_id == ctx.execution_id && node_id == ctx.node_id => { + let primary = if decision.is_empty() { + decisions.first().cloned().unwrap_or_default() + } else { + decision + }; + return Ok(NodeOutput::from_value(json!({ "decision": primary }))); + } + Ok(_) => continue, + Err(_) => return Err(anyhow::anyhow!("事件总线关闭,审批无法完成")), + }, + _ = tokio::time::sleep_until(deadline) => { + return Err(anyhow::anyhow!("人工审批等待超时({timeout_secs}s)")); + } + } + } + } + + fn schema(&self) -> NodeSchema { + NodeSchema { + params: json!(null), + output: json!(null), + } + } + + fn is_blocking(&self) -> bool { + true + } + + fn node_type(&self) -> &str { + "human" + } +} + +/// 从事件流中等待一条 HumanApprovalRequest(跳过 NodeStarted/NodeCompleted 等其他事件)。 +async fn wait_approval_request(mut rx: EventSubscriber) -> WorkflowEvent { + tokio::time::timeout(Duration::from_millis(1000), async move { + loop { + if let Ok(ev @ WorkflowEvent::HumanApprovalRequest { .. }) = rx.recv().await { + return ev; + } + } + }) + .await + .expect("应收到 HumanApprovalRequest,实际超时") +} + +/// 端到端主链路:前驱 sleep(a) → human(b)。 +/// 执行到 human 后挂起等审批 → 外部 approve(模拟 approve_human_approval IPC) → 返回结果、工作流完成。 +#[tokio::test] +async fn human_workflow_end_to_end_approval() { + let bus = EventBus::new(); + let rx = bus.subscribe(); // 先 subscribe 再执行(broadcast 不回放) + let exec_id = "exec-e2e-approve"; + + // DagDef(可序列化定义)→ NodeRegistry.build_dag → 运行时 Dag + let mut def = DagDef::new(); + def.add_node("a", "sleep", json!({})); + def.add_node("b", "human", json!({})); + def.add_edge("a", "b"); + let mut registry = NodeRegistry::new(); + registry.register("sleep", |_| Box::new(SleepNode { sleep_ms: 10 })); + registry.register("human", |_| Box::new(ApprovalNode)); + let dag = registry.build_dag(&def).expect("build_dag 应成功"); + + let mut executor = DagExecutor::new(bus.clone(), exec_id.into()); + let sm = executor.state_machine(); // 共享状态机(spawn 后仍可读) + let run_handle = tokio::spawn(async move { executor.run(&dag, json!({})).await }); + + // a 层完成、b 层 human subscribe + send Request 后应收到审批请求 + let request = wait_approval_request(rx).await; + match request { + WorkflowEvent::HumanApprovalRequest { + node_id, title, options, .. + } => { + assert_eq!(node_id, "b"); + assert_eq!(title, "确认发布"); + assert_eq!(options, vec!["同意".to_string(), "拒绝".to_string()]); + } + other => panic!("期望 HumanApprovalRequest,收到 {:?}", other), + } + + // 外部 approve(模拟 approve_human_approval IPC 发送 HumanApprovalResponse 到总线) + bus.send(WorkflowEvent::HumanApprovalResponse { + execution_id: exec_id.into(), + node_id: "b".to_string(), + decision: "同意".into(), + decisions: vec![], + comment: Some("可以发布".into()), + }) + .await; + + let outputs = run_handle + .await + .expect("run 不应 panic") + .expect("run 应成功"); + assert!(outputs.contains_key("a"), "outputs 应含前驱节点 a"); + assert_eq!(outputs["b"].data["decision"], json!("同意")); + assert_eq!(sm.get(&"a".to_string()), NodeStatus::Completed); + assert_eq!(sm.get(&"b".to_string()), NodeStatus::Completed); +} + +/// 取消路径:外部 set_cancelled(模拟 cancel_workflow_node IPC 写共享状态机) +/// → human select! 轮询 is_cancelled → Err → 状态保持 Cancelled。 +#[tokio::test] +async fn human_workflow_cancel_via_shared_state_machine() { + let bus = EventBus::new(); + let rx = bus.subscribe(); + let exec_id = "exec-e2e-cancel"; + + let mut def = DagDef::new(); + def.add_node("h", "human", json!({})); + let mut registry = NodeRegistry::new(); + registry.register("human", |_| Box::new(ApprovalNode)); + let dag = registry.build_dag(&def).expect("build_dag 应成功"); + + let mut executor = DagExecutor::new(bus.clone(), exec_id.into()); + let sm = executor.state_machine(); + let run_handle = tokio::spawn(async move { executor.run(&dag, json!({})).await }); + + // 等 human 发出审批请求(确认已 subscribe + send + 进入 select! 轮询) + let _ = wait_approval_request(rx).await; + + // 外部 cancel:写共享状态机置 Cancelled,humar select! 的 cancel_tick 会读到 + sm.set_cancelled("h".to_string()); + + let result = run_handle.await.unwrap(); + assert!(result.is_err(), "取消应致 run 返回 Err"); + let err = result.unwrap_err().to_string(); + assert!( + err.contains("取消") || err.contains("执行失败"), + "应返回取消相关错误,实际: {err}" + ); + // 状态保持 Cancelled(未被 set_failed 覆盖为 Failed)—— 取消终态不被失败语义覆盖 + assert_eq!(sm.get(&"h".to_string()), NodeStatus::Cancelled); +} + +/// 挂起等待路径:收到请求后不 approve → human 挂起直到超时 → run 返 Err。 +/// 证明节点真在等待外部响应(而非空转直接返回)。 +#[tokio::test] +async fn human_workflow_waits_then_times_out_without_approve() { + let bus = EventBus::new(); + let rx = bus.subscribe(); + let exec_id = "exec-e2e-timeout"; + + let mut def = DagDef::new(); + def.add_node("h", "human", json!({ "timeout_secs": 1 })); + let mut registry = NodeRegistry::new(); + registry.register("human", |_| Box::new(ApprovalNode)); + let dag = registry.build_dag(&def).expect("build_dag 应成功"); + + let mut executor = DagExecutor::new(bus.clone(), exec_id.into()); + let run_handle = tokio::spawn(async move { executor.run(&dag, json!({})).await }); + + let _ = wait_approval_request(rx).await; + let result = run_handle.await.unwrap(); + assert!(result.is_err(), "无审批响应应超时失败"); + // executor 包了「节点 h 执行失败」context,用 {:#} 展开整条错误链断言根因 + let chain = format!("{:#}", result.unwrap_err()); + assert!(chain.contains("超时"), "应返回超时错误,实际: {chain}"); +} diff --git a/src-tauri/src/commands/ai/agentic/mod.rs b/src-tauri/src/commands/ai/agentic/mod.rs index 1900de9..b35e883 100644 --- a/src-tauri/src/commands/ai/agentic/mod.rs +++ b/src-tauri/src/commands/ai/agentic/mod.rs @@ -151,6 +151,23 @@ pub const GOAL_MAX_CHARS: usize = 500; /// 防无限膨胀:每次发消息提取的目标追加到 pinned_goals vec,超上限时淘汰最早目标。 pub const MAX_GOALS: usize = 5; +/// 技能清单注入开关:system_prompt 是否追加本机技能清单(名称+描述),让 LLM 感知可用技能(默认 true)。 +/// +/// 机制化:LLM 只读清单,用户请求命中某技能用途时能引导触发(/技能名),不做复杂执行逻辑。 +/// 清单来自 skills.rs 进程内缓存(skills/commands/plugins 三类去重)。false(回退)→ 整块跳过, +/// system_prompt 零变化,退回改动前行为(单点回退,不影响 augmentation 的 / 联想)。 +pub const SKILL_LIST_INJECT_ENABLED: bool = true; + +/// 技能清单注入最大条数(默认 20)。 +/// +/// 防本机大量技能全量进 system_prompt 占 sys_tokens 预算;取前 N 条(name 去重序),截断描述。 +pub const SKILL_LIST_MAX_ITEMS: usize = 20; + +/// 技能清单单条描述截断长度(默认 80 字符)。 +/// +/// 清单只做提示,截断防长描述(部分技能 description 数百字)撑爆 prompt,截断后接省略号。 +pub const SKILL_DESC_MAX_CHARS: usize = 80; + /// G4 目标感知降级:话题标记 insert 当 pinned_goals 存在时跳过(默认 true)。 /// /// 治 R2(话题标记反向误导):G1 目标钉扎生效后每轮 system_prompt 已含目标,topic marker 的 @@ -1044,6 +1061,45 @@ pub(crate) async fn run_agentic_loop( ); system_prompt = format!("{}\n\n{}\n\n{}", system_prompt, env_prompt, behavior_prompt); + // 技能清单注入:把本机扫描到的技能(名称+描述)追加到 system_prompt,让 LLM 感知可用技能。 + // + // 治「AI 不知道有技能可引导」:skills.rs 联想已做(前端 / 浮层),但 AI 本身对技能无感知, + // 用户需求命中某技能(如"发布flux")时 AI 不会引导。机制化——只注入清单(不执行技能逻辑, + // 执行仍由用户在 Claude Code 侧 / 技能名 触发),LLM 据此提示用户可用技能。 + // 清单来自 skills_cached() 进程内缓存(懒扫描一次,后续零开销);截断条数+描述防占预算。 + if SKILL_LIST_INJECT_ENABLED { + let skills = crate::commands::ai::skills::skills_cached().await; + if !skills.is_empty() { + let lines: Vec = skills + .iter() + .take(SKILL_LIST_MAX_ITEMS) + .map(|s| { + let desc_flat = s.description.replace('\n', " "); + let desc: String = desc_flat.chars().take(SKILL_DESC_MAX_CHARS).collect(); + let desc = if desc_flat.chars().count() > SKILL_DESC_MAX_CHARS { + format!("{}…", desc) + } else { + desc + }; + format!("- {}: {}", s.name, desc) + }) + .collect(); + if !lines.is_empty() { + system_prompt = format!( + "{}\n\n## 可用技能\n以下为本机 Claude 技能/命令(名称: 描述)。用户请求命中某技能用途时,提示用户输入 /技能名 触发,或引导其按需使用:\n{}", + system_prompt, + lines.join("\n") + ); + tracing::debug!( + conv_id = %conv_id, + count = lines.len(), + "[ai] 技能清单注入:已追加 {} 条技能到 system_prompt", + lines.len() + ); + } + } + } + // T4: 工作流 DAG 注入 — 当会话关联工作流时,将活跃路径注入 system prompt { let session = session_arc.lock().await; diff --git a/src-tauri/src/commands/ai/tool_registry.rs b/src-tauri/src/commands/ai/tool_registry.rs index e3c067e..2542e80 100644 --- a/src-tauri/src/commands/ai/tool_registry.rs +++ b/src-tauri/src/commands/ai/tool_registry.rs @@ -1365,15 +1365,18 @@ mod tests { // list_module_dependencies/add_module_dependency/remove_module_dependency,扩展 tools/task_graph.rs, // 复用 module.rs CRUD 逻辑;delete_module 级联清理依赖边防外键约束)。 // 72 = 42 data + 14 file + 1 http + 1 fetch_url + 1 fetch_search + 1 generate_image + 1 get_app_config + 11 local_proxy。 + // diff_files(2026-08-08): file 层 14→15(新增任意两文件内容对比,只读 Low, + // 补 git_diff 仅仓库内缺口;复用 generate_diff LCS unified diff;namespace ALWAYS_LARGE_TOOLS 既有引用生效)。 + // 73 = 42 data + 15 file + 1 http + 1 fetch_url + 1 fetch_search + 1 generate_image + 1 get_app_config + 11 local_proxy。 assert_eq!( registry.len(), - 72, - "工具总数应为 72(42 data + 14 file + 1 http + 1 fetch_url + 1 fetch_search + 1 generate_image + 1 get_app_config + 11 local_proxy),实际 {}", registry.len() + 73, + "工具总数应为 73(42 data + 15 file + 1 http + 1 fetch_url + 1 fetch_search + 1 generate_image + 1 get_app_config + 11 local_proxy),实际 {}", registry.len() ); // 工具名集合基线:防 rename / 漏注册 / 误删除。 // data 层 42 个(持 db):CRUD/状态机/工作流/知识图谱任务关联 + 项目事件流 + 基础设施配置 + git 工具 + 工程模块写工具 - // file 层 14 个(不持 db):命令/读/列/写/改/元/追加/删/移/搜/grep/环境探测/符号解析/下载 + // file 层 15 个(不持 db):命令/读/列/写/改/元/追加/删/移/搜/grep/环境探测/符号解析/下载/对比 // http 层 1 个(不持 db):http_request let mut expected: Vec<&str> = vec![ // ── data 层 (42) ── @@ -1400,16 +1403,18 @@ mod tests { "git_status", "git_diff", "git_log", // Git 写工具 "git_commit", "git_branch", "git_merge", - // ── file 层 (14) ──(run_command 注册顺序已移至末位降低 LLM 偏好, + // ── file 层 (15) ──(run_command 注册顺序已移至末位降低 LLM 偏好, // 集合断言经 sort 后与顺序无关,仅守护工具名不漂移。grep 新增 F-260621; // detect_environment 新增 L1 环境感知 设计 §2.1;read_symbol 新增 AST 代码智能; - // download_file 新增跨平台 URL→文件流式下载) + // download_file 新增跨平台 URL→文件流式下载;diff_files 新增任意两文件对比) "read_file", "read_symbol", "list_directory", "write_file", "patch_file", "file_info", "append_file", "delete_file", "rename_file", "search_files", "run_command", "grep", "detect_environment", // download_file(URL→文件流式下载,写文件需 allowed_dirs,注册在 register_file_tools) "download_file", + // diff_files(任意两文件内容对比,只读 Low,补 git_diff 仅仓库内缺口) + "diff_files", // ── http 层 (1) ── "http_request", // ── fetch_url 层 (1) ──(URL → markdown 文档嗅探,只读 GET,与 http_request 分工) diff --git a/src-tauri/src/commands/ai/tools/file.rs b/src-tauri/src/commands/ai/tools/file.rs index 84ad98c..7b51d0c 100644 --- a/src-tauri/src/commands/ai/tools/file.rs +++ b/src-tauri/src/commands/ai/tools/file.rs @@ -259,6 +259,46 @@ pub fn register( } ); + // ── diff_files (Low 只读, 任意两文件对比) ── + // 补 git_diff 缺口:git_diff 只比仓库内版本,本工具对比任意两文件(绝对路径), + // 复用 generate_diff(LCS unified diff)生成差异,只读无副作用。 + declare_tool!( + registry, + allowed_dirs: Arc>, + "diff_files", + "对比任意两个文件的差异,返回 unified diff(+/- 行前缀,基于 LCS)。参数 path_a/path_b 为两文件绝对路径,均须在授权目录内。只读无副作用,用于两文件/两版本内容对比(git_diff 仅仓库内,本工具补任意两文件缺口)", + RiskLevel::Low, + schema: object_schema(vec![("path_a", "string", true), ("path_b", "string", true)]), + args => { + let snap = allowed_dirs.read().await.clone(); + let a_resolved = resolve_workspace_path_with_allowed( + args["path_a"].as_str().ok_or_else(|| anyhow::anyhow!("缺少 path_a 参数"))?, + &snap, + )?; + let b_resolved = resolve_workspace_path_with_allowed( + args["path_b"].as_str().ok_or_else(|| anyhow::anyhow!("缺少 path_b 参数"))?, + &snap, + )?; + let a_path = a_resolved.to_str().ok_or_else(|| anyhow::anyhow!("path_a 含非法字符"))?.to_string(); + let b_path = b_resolved.to_str().ok_or_else(|| anyhow::anyhow!("path_b 含非法字符"))?.to_string(); + if a_path == b_path { + anyhow::bail!("两个路径相同,无需对比: {}", a_path); + } + // 读取两文件文本(1MB 上限 + 二进制拦截,解码对齐 read_file 的 UTF-16 BOM 优先) + let (content_a, size_a) = read_text_file(&a_path).await?; + let (content_b, size_b) = read_text_file(&b_path).await?; + let diff = generate_diff(&content_a, &content_b); + Ok(serde_json::json!({ + "path_a": a_path, + "path_b": b_path, + "size_a": size_a, + "size_b": size_b, + "identical": content_a == content_b, + "diff": diff, + })) + } + ); + // ── list_directory (Low 只读) ── declare_tool!( registry, @@ -1232,6 +1272,38 @@ pub fn register( ); } +// ============================================================ +// diff_files 辅助 — 读取文本文件内容(两文件对比共用) +// ============================================================ + +/// 读取文本文件内容,返回 (内容, 字节数)。 +/// +/// diff_files 两文件对比共用:1MB 上限 + 二进制拦截 + UTF-16 BOM 优先解码 +/// (复用 decode_bytes_to_string,对齐 read_file/patch_file 的读取语义,单真相源)。 +async fn read_text_file(path: &str) -> anyhow::Result<(String, u64)> { + use tokio::fs::File; + use tokio::io::AsyncReadExt; + let mut file = File::open(path).await + .map_err(|e| { + if e.kind() == std::io::ErrorKind::NotFound { + anyhow::anyhow!("文件不存在: {}", path) + } else { + anyhow::anyhow!("无法访问文件 {}: {}", path, e) + } + })?; + let metadata = file.metadata().await + .map_err(|e| anyhow::anyhow!("读取元数据失败 {}: {}", path, e))?; + if metadata.len() > 1_048_576 { + anyhow::bail!("文件超过 1MB 限制 ({} 字节)", metadata.len()); + } + let mut raw_bytes = Vec::with_capacity(metadata.len() as usize); + file.read_to_end(&mut raw_bytes).await + .map_err(|e| anyhow::anyhow!("读取文件失败: {}", e))?; + let content = decode_bytes_to_string(&raw_bytes) + .map_err(|_| anyhow::anyhow!("文件非 UTF-8 文本(疑似二进制),无法对比: {}", path))?; + Ok((content, metadata.len())) +} + // ============================================================ // patch_file 失败提示辅助 — old_text 不匹配时给 LLM 相近锚点 // ============================================================ diff --git a/src/api/types.ts b/src/api/types.ts index 8afeb77..b0b4ef7 100644 --- a/src/api/types.ts +++ b/src/api/types.ts @@ -288,6 +288,9 @@ export interface WorkflowEventPayload { } } +/** 工作流 DAG 节点实时状态(WorkflowDagDisplay 高亮 + workflow store 共享推导) */ +export type NodeStatus = 'pending' | 'running' | 'completed' | 'failed' | 'skipped' + // ============================================================ // 通用 // ============================================================ diff --git a/src/components/workflow/WorkflowDagDisplay.vue b/src/components/workflow/WorkflowDagDisplay.vue index c3282e6..5dfcc89 100644 --- a/src/components/workflow/WorkflowDagDisplay.vue +++ b/src/components/workflow/WorkflowDagDisplay.vue @@ -65,6 +65,8 @@ @@ -1205,6 +1240,20 @@ onUnmounted(() => { } .workflow-log-panel { flex: 0 0 auto; } +/* B-41 工作流实时进度面板(与 TaskDetail wf-progress 同款视觉) */ +.wf-progress-panel { flex: 0 0 auto; } +.wf-progress { + display: flex; + flex-direction: column; + gap: 2px; + margin-top: 6px; + font-size: 12px; + color: var(--df-text-secondary); +} +.wf-progress-count { color: var(--df-text-dim); } +.wf-progress-hint { color: var(--df-success); } +.wf-progress-hint-fail { color: var(--df-danger); } + /* ===== 面板 ===== */ /* .panel / .panel-header 基础样式已收敛至全局 components.css(DRY 收口 B-260619), 此处仅保留本组件特有 .task-count。 */ diff --git a/src/views/TaskDetail.vue b/src/views/TaskDetail.vue index e0bffdd..2b6626b 100644 --- a/src/views/TaskDetail.vue +++ b/src/views/TaskDetail.vue @@ -237,6 +237,7 @@ import { useRoute, useRouter } from 'vue-router' import { listen } from '@tauri-apps/api/event' import { taskApi, projectApi, ideaApi } from '@/api' import { workflowApi } from '@/api' +import { useProjectStore } from '@/stores/project' import { formatDate } from '@/utils/time' import { useRendered } from '@/composables/useMarkdown' import { @@ -245,14 +246,15 @@ import { priorityLabel, priorityClass, } from '../constants/project' -import type { TaskRecord, ProjectRecord, IdeaRecord, DfDataChangedPayload, WorkflowEventPayload } from '@/api/types' +import type { TaskRecord, ProjectRecord, IdeaRecord, DfDataChangedPayload } from '@/api/types' import TaskOutputCard from '@/components/task/TaskOutputCard.vue' import WorkflowDagDisplay from '@/components/workflow/WorkflowDagDisplay.vue' -import type { NodeStatus } from '@/components/workflow/WorkflowDagDisplay.vue' const { t } = useI18n() const route = useRoute() const router = useRouter() +// B-41: 共享 workflow store(liveEvents 全局收集 + workflowProgress 派生进度) +const store = useProjectStore() const loading = ref(true) const errorMsg = ref('') @@ -386,32 +388,47 @@ async function confirmSubtask() { // F-260616-06 ①-1 / B-41 工作流推进状态(与手动 advance 的 advancing 独立,互不干扰) // ------------------------------------------------------------ // wfAdvancing:推进中态;wfExecId:当前监听的工作流执行 ID(过滤 workflow-event); -// wfTotalNodes/wfDoneCount/wfRunningNode:轻量进度(听 NodeStarted/NodeCompleted); +// wfTotalNodes/wfDoneCount/wfRunningNode:轻量进度(由共享 store 按 execId 派生); // wfResult: 'completed' | 'failed' | null —— 终态提示,完成/失败后自动清空(由新一次推进重置)。 const wfAdvancing = ref(false) const wfExecId = ref(null) -const wfTotalNodes = ref(0) -const wfDoneCount = ref(0) -const wfRunningNode = ref('') const wfResult = ref<'completed' | 'failed' | null>(null) // SW-260618-21: 终态提示 timer 引用,卸载时清理防写已销毁 ref(对齐 Projects.vue _toastTimer 模式) let _wfResultTimer: ReturnType | null = null +// B-41: 共享工作流进度推导。workflow store 已全局收集 liveEvents,此处按 execId +// 派生节点状态/进度,与 ProjectDetail 同源(store.workflowProgress),替代原本地 onEvent 逐事件累加。 +const wfProgress = computed(() => store.workflowProgress(wfExecId.value)) +const wfNodeStatuses = computed(() => wfProgress.value.nodeStatuses) +const wfDoneCount = computed(() => wfProgress.value.doneCount) +const wfTotalNodes = computed(() => wfProgress.value.totalNodes) +const wfRunningNode = computed(() => wfProgress.value.runningNode) +const wfDoneTotal = computed(() => wfTotalNodes.value) + const wfProgressHint = computed(() => wfResult.value !== null) const wfCompletedHint = computed(() => wfResult.value === 'completed') const wfFailedHint = computed(() => wfResult.value === 'failed') -// 工作流 DAG 结构展示 +// 终态 → 瞬态提示:派生 result 置位时收 wfAdvancing + 3s 后清 wfResult(仅 UI hint,数据不改) +watch(() => wfProgress.value.result, (r) => { + if (!r) return + wfResult.value = r + wfAdvancing.value = false + if (_wfResultTimer) clearTimeout(_wfResultTimer) + _wfResultTimer = setTimeout(() => { + if (wfResult.value === r) wfResult.value = null + _wfResultTimer = null + }, 3000) +}) + +// 工作流 DAG 结构展示(执行记录读回,节点状态由共享派生) const wfDagJson = ref('') -const wfNodeStatuses = ref>({}) async function refreshWorkflowDag(execId: string) { - wfNodeStatuses.value = {} try { const record = await workflowApi.getExecution(execId) if (record?.dag_json) wfDagJson.value = record.dag_json } catch { /* 静默 */ } } -const wfDoneTotal = computed(() => wfTotalNodes.value) const taskId = computed(() => route.params.id as string) @@ -568,9 +585,6 @@ async function handleWorkflowAdvance(target: string) { if (!task.value || wfAdvancing.value || advancing.value) return wfAdvancing.value = true wfResult.value = null - wfRunningNode.value = '' - wfDoneCount.value = 0 - wfTotalNodes.value = 0 try { const execId = await workflowApi.run( `task-advance:${target}`, @@ -589,64 +603,8 @@ async function handleWorkflowAdvance(target: string) { } } -// B-41: 处理单个工作流事件(由 onEvent 回调 dispatch)。按 execution_id 过滤当前推进。 -function handleWorkflowEvent(payload: WorkflowEventPayload) { - if (!wfExecId.value || payload.execution_id !== wfExecId.value) return - const evt = payload.event - switch (evt?.type) { - case 'node_started': { - // 第一次 NodeStarted 时初始化 total(后端模板节点数,无法前端预知, - // 用「已启动 + 已完成」近似 total,显示已完成/已启动进度) - const node = String(evt.node_id ?? '') - wfRunningNode.value = node - wfNodeStatuses.value = { ...wfNodeStatuses.value, [node]: 'running' } - if (wfTotalNodes.value === 0) wfTotalNodes.value = 1 - else wfTotalNodes.value = Math.max(wfTotalNodes.value, wfDoneCount.value + 1) - break - } - case 'node_completed': { - wfDoneCount.value += 1 - if (wfRunningNode.value) { - wfNodeStatuses.value = { ...wfNodeStatuses.value, [wfRunningNode.value]: 'completed' } - } - wfRunningNode.value = '' - break - } - case 'workflow_completed': { - wfAdvancing.value = false - wfRunningNode.value = '' - wfResult.value = 'completed' - // 终态提示保留 3s 后清空(任务本体由 df-data-changed 刷新,提示不影响数据) - if (_wfResultTimer) clearTimeout(_wfResultTimer) - _wfResultTimer = setTimeout(() => { - if (wfResult.value === 'completed') wfResult.value = null - _wfResultTimer = null - }, 3000) - break - } - case 'workflow_failed': - case 'node_failed': { - if (evt?.node_id) { - const node = String(evt.node_id) - wfNodeStatuses.value = { ...wfNodeStatuses.value, [node]: 'failed' } - } - if (evt?.type === 'workflow_failed') { - wfAdvancing.value = false - wfRunningNode.value = '' - wfResult.value = 'failed' - if (_wfResultTimer) clearTimeout(_wfResultTimer) - _wfResultTimer = setTimeout(() => { - if (wfResult.value === 'failed') wfResult.value = null - _wfResultTimer = null - }, 3000) - } - // node_failed: 单节点失败,工作流后续会发 workflow_failed,此处不提前改终态 - break - } - default: - break - } -} +// B-41: 进度状态已由共享 store 按 execId 派生(store.workflowProgress),不再本地逐事件累加。 +// 原有 handleWorkflowEvent 已移除 —— 事件仍由 workflow store 全局收集,TaskDetail 只读派生结果。 async function load() { loading.value = true @@ -689,8 +647,6 @@ watch(renderedDesc, () => { measureDescHeight() }) // 不享受 store 全局 df-data-changed 监听(该监听只刷 store.tasks 列表,不含本视图的当前 task 单体), // 故此本地监听 entity=task/project 时重载当前 task。 let _unlistenDataChanged: (() => void) | null = null -// B-41: 工作流推进事件 unlistener(听 NodeStarted/NodeCompleted/WorkflowFailed 显示轻量进度) -let _unlistenWorkflowEvent: (() => void) | null = null onMounted(async () => { ensureLoaded() // 后台预热 Markdown 渲染器(模块单例,与 AiChat 共享),不阻塞 load @@ -709,18 +665,17 @@ onMounted(async () => { } catch (e) { console.error('[TaskDetail] 启动 df-data-changed 监听失败:', e) } - // B-41: 工作流推进事件监听(复用 workflowApi.onEvent, 与 workflow store 的全局监听并存无冲突 — - // 多个 listen 各自接收全量事件,本视图按 execution_id 过滤只处理自己发起的推进) + // B-41: 工作流进度由 workflow store 全局 liveEvents 推导(store.workflowProgress), + // 此处只需确保 store 事件监听已启动(幂等,共享单例;按 execution_id 过滤由派生层完成)。 try { - _unlistenWorkflowEvent = await workflowApi.onEvent(handleWorkflowEvent) + await store.startEventListener() } catch (e) { - console.error('[TaskDetail] 启动 workflow-event 监听失败:', e) + console.error('[TaskDetail] 启动 workflow 事件监听失败:', e) } }) onBeforeUnmount(() => { if (_unlistenDataChanged) { _unlistenDataChanged(); _unlistenDataChanged = null } - if (_unlistenWorkflowEvent) { _unlistenWorkflowEvent(); _unlistenWorkflowEvent = null } // F-260805:移除子任务快捷菜单关闭监听 document.removeEventListener('click', closeChildMenu) // SW-260618-21: 清终态提示 timer 防卸载后写已销毁 ref