Browse Source

1. 实现了 RAG AI Bridge 进程管理器,由 Java 端通过 ProcessBuilder 启动 Python 服务并以环境变量注入 LLM 配置,替代原有独立进程启动方式,统一了模型配置的下发链路;
2. 修复了 Embedding Bridge 的 Milvus 标量结果丢失和集合 released 状态问题,新增 _ensure_loaded 与 _run_read 方法保证查询前集合已加载并统一走读取通道;
3. 实现了工作流同层节点并行执行,新增独立 nodeExecutor 线程池避免与主线程池争用导致死锁,SseEmitter 推送加 synchronized 保证多 worker 并发下的线程安全;
4. 实现了工作流上下文系统提示注入,新增 ContextPromptHelper 统一格式化前置节点变量为系统提示,覆盖 LLM、智能操作、技能、Hermes Agent、Hermes 智能操作五类执行器,解决下游节点看不到上游输出的问题;
5. 将工作流最大执行超时改为无限制(Long.MAX_VALUE),避免长任务被调度器强制中断;
6. 实现了 RAG 治理页“应用审核后的授权配置”自动建立数据源绑定,新增按 (knowledgeBaseId, sourceType, sourceId) 的 upsert 接口,解决用户未预先创建绑定记录导致报错的问题;
7. 修复了图 RAG 修复接口超时问题,/repair 端点由硬编码 20 秒改为复用 SETTINGS.llm.timeout(默认 60 秒),与 text2cypher 主链路对齐;
8. 增强了 Text2Cypher 与 Repair Prompt 的 Schema 约束,强制每个节点模式必须显式声明一个来自 allowedLabels 的标签,禁止凭空创造、翻译或复数化标签名,修复 bare node 与 nonexistent label 两类校验失败;
9. 修复了 RAG AI Bridge 脚本路径重复问题,scriptPath 默认值由 `backend/rag-ai-bridge/server.py` 改为 `rag-ai-bridge/server.py`,避免与 user.dir 拼接后出现 `backend/backend/` 双层路径;
10. 新增了知识检索工作流节点,支持在 DAG 中通过 KnowledgeRetrievalNode 调用 RAG 检索能力并接入工作流变量;
11. 更新了 prompt.md 需求记录与 application.yml.example 配置示例。

weisijie 1 tháng trước cách đây
mục cha
commit
94a9f24a5d
26 tập tin đã thay đổi với 1341 bổ sung136 xóa
  1. 1 0
      .gitignore
  2. 9 3
      AGENTS.md
  3. 118 61
      backend/embedding-bridge/backend/indexing/milvus_client.py
  4. 8 1
      backend/rag-ai-bridge/server.py
  5. 39 2
      backend/src/main/java/com/agent/management/config/RagAiBridgeProperties.java
  6. 83 0
      backend/src/main/java/com/agent/management/engine/ContextPromptHelper.java
  7. 2 2
      backend/src/main/java/com/agent/management/engine/WorkflowEngine.java
  8. 126 46
      backend/src/main/java/com/agent/management/engine/WorkflowLevelExecutor.java
  9. 3 1
      backend/src/main/java/com/agent/management/engine/executor/AgentExecutor.java
  10. 2 1
      backend/src/main/java/com/agent/management/engine/executor/HermesAgentExecutor.java
  11. 4 1
      backend/src/main/java/com/agent/management/engine/executor/HermesSmartActionExecutor.java
  12. 257 0
      backend/src/main/java/com/agent/management/engine/executor/KnowledgeRetrievalExecutor.java
  13. 4 1
      backend/src/main/java/com/agent/management/engine/executor/LlmExecutor.java
  14. 7 4
      backend/src/main/java/com/agent/management/engine/executor/SmartActionExecutor.java
  15. 1 1
      backend/src/main/java/com/agent/management/rag/bridge/RagAiBridgeClient.java
  16. 211 0
      backend/src/main/java/com/agent/management/rag/bridge/RagAiBridgeProcessManager.java
  17. 2 2
      backend/src/main/java/com/agent/management/rag/controller/KnowledgeBaseRagDebugController.java
  18. 35 2
      backend/src/main/java/com/agent/management/rag/kb/RagKnowledgeBaseConfigService.java
  19. 2 2
      backend/src/main/java/com/agent/management/repository/RagKnowledgeSourceBindingRepository.java
  20. 7 1
      backend/src/main/resources/application.yml.example
  21. 8 0
      frontend/src/api/rag.js
  22. 69 0
      frontend/src/components/workflow/nodes/KnowledgeRetrievalNode.vue
  23. 15 0
      frontend/src/utils/ioInference.js
  24. 3 3
      frontend/src/views/knowledge/RagGovernance.vue
  25. 241 2
      frontend/src/views/workflow/WorkflowEditor.vue
  26. 84 0
      prompt.md

+ 1 - 0
.gitignore

@@ -30,6 +30,7 @@ backend/src/main/resources/application.yml
 !.env.example
 backend/rag-ai-bridge/.env
 backend/rag-ai-bridge/config.yaml
+docker-compose.yml
 
 # ===== 日志 =====
 logs/

+ 9 - 3
AGENTS.md

@@ -1,6 +1,6 @@
-始终使用简体中文回复
+# 通用准则
 
----
+**始终使用简体中文回复**。
 
 **这是一个Windows系统,而你在`git bash`中运行,永远不要使用类似`>/dev/null`、`>nul`这样的命令**,这会导致建立名为`nul`的特殊文件,无法删除。
 
@@ -8,7 +8,11 @@
 
 **永远不要执行“杀死所有Java进程、杀死所有Node进程”这样的操作**,请务必根据端口号或文件名精准到筛选出特定进程。
 
----
+**针对比较复杂的网页操作,不要使用`playwright`来进行测试**,既慢又耗费 token,交给用户进行手动测试。
+
+除非我明确要求,否则**不要帮我启动前后端服务**。如果为了测试需要启动,测试完成后,关闭服务,由用户手动启动。
+
+# 编程行为准则
 
 行为准则,旨在减少常见的 LLM 编码错误。可根据项目特定指令按需合并。
 
@@ -76,4 +80,6 @@
 
 ---
 
+# Token 节省准则
+
 @RTK.md

+ 118 - 61
backend/embedding-bridge/backend/indexing/milvus_client.py

@@ -64,6 +64,26 @@ class MilvusStore:
         with milvus_client_session(self._settings) as client:
             return operation(client)
 
+    @staticmethod
+    def _ensure_loaded(client: MilvusClient, collection_name: str) -> None:
+        """幂等加载 collection。
+
+        Milvus 服务重启后 collection 默认 released,直接 search/query 会抛
+        'call load() before search'。已 loaded 的 collection 重复调用会被
+        服务端识别为无操作。collection 不存在 / 正在加载等异常忽略,真正
+        的状态错误会在后续操作时抛出明确信息。
+        """
+        try:
+            client.load_collection(collection_name)
+        except Exception:
+            pass
+
+    def _run_read(self, operation: Callable[[MilvusClient], T]) -> T:
+        """读取类操作统一入口:先幂等 load collection,再执行。"""
+        with milvus_client_session(self._settings) as client:
+            MilvusStore._ensure_loaded(client, self.collection_name)
+            return operation(client)
+
     @contextmanager
     def session(self) -> Iterator[MilvusClient]:
         """同一业务流(如整次上传)内复用一条连接,用毕即关。"""
@@ -72,65 +92,102 @@ class MilvusStore:
 
     @staticmethod
     def ensure_collection(client: MilvusClient, collection_name: str, dense_dim: int) -> None:
-        if client.has_collection(collection_name):
-            return
-
-        schema = client.create_schema(auto_id=True, enable_dynamic_field=True)
-        schema.add_field("id", DataType.INT64, is_primary=True, auto_id=True)
-        schema.add_field("dense_embedding", DataType.FLOAT_VECTOR, dim=dense_dim)
-        schema.add_field("sparse_embedding", DataType.SPARSE_FLOAT_VECTOR)
-        schema.add_field(
-            "text",
-            DataType.VARCHAR,
-            max_length=65535,
-            enable_analyzer=True,
-            analyzer_params={"type": "standard"},
-            enable_match=True,
-        )
-        schema.add_field("document_id", DataType.INT64)
-        schema.add_field("filename", DataType.VARCHAR, max_length=255)
-        schema.add_field("file_type", DataType.VARCHAR, max_length=50)
-        schema.add_field("file_path", DataType.VARCHAR, max_length=1024)
-        schema.add_field("page_number", DataType.INT64)
-        schema.add_field("chunk_idx", DataType.INT64)
-        schema.add_field("chunk_id", DataType.VARCHAR, max_length=512)
-        schema.add_field("parent_chunk_id", DataType.VARCHAR, max_length=512)
-        schema.add_field("root_chunk_id", DataType.VARCHAR, max_length=512)
-        schema.add_field("chunk_level", DataType.INT64)
-
-        bm25_function = Function(
-            name="text_bm25_emb",
-            function_type=FunctionType.BM25,
-            input_field_names=["text"],
-            output_field_names=["sparse_embedding"],
-        )
-        schema.add_function(bm25_function)
-
-        index_params = client.prepare_index_params()
-        index_params.add_index(
-            field_name="dense_embedding",
-            index_type="HNSW",
-            metric_type="IP",
-            params={"M": 16, "efConstruction": 256},
-        )
-        index_params.add_index(
-            field_name="sparse_embedding",
-            index_type="SPARSE_INVERTED_INDEX",
-            metric_type="BM25",
-            params={"drop_ratio_build": 0.2},
-        )
+        # 不论 collection 是否已存在,最终都要做索引对账:
+        # - 历史已存在的 collection 可能因旧版异常处理被静默吞掉索引创建失败
+        # - 新建过程中若 create_collection 部分成功(collection 已建、索引未建全)也需要补救
+        if not client.has_collection(collection_name):
+            schema = client.create_schema(auto_id=True, enable_dynamic_field=True)
+            schema.add_field("id", DataType.INT64, is_primary=True, auto_id=True)
+            schema.add_field("dense_embedding", DataType.FLOAT_VECTOR, dim=dense_dim)
+            schema.add_field("sparse_embedding", DataType.SPARSE_FLOAT_VECTOR)
+            schema.add_field(
+                "text",
+                DataType.VARCHAR,
+                max_length=65535,
+                enable_analyzer=True,
+                analyzer_params={"type": "standard"},
+                enable_match=True,
+            )
+            schema.add_field("document_id", DataType.INT64)
+            schema.add_field("filename", DataType.VARCHAR, max_length=255)
+            schema.add_field("file_type", DataType.VARCHAR, max_length=50)
+            schema.add_field("file_path", DataType.VARCHAR, max_length=1024)
+            schema.add_field("page_number", DataType.INT64)
+            schema.add_field("chunk_idx", DataType.INT64)
+            schema.add_field("chunk_id", DataType.VARCHAR, max_length=512)
+            schema.add_field("parent_chunk_id", DataType.VARCHAR, max_length=512)
+            schema.add_field("root_chunk_id", DataType.VARCHAR, max_length=512)
+            schema.add_field("chunk_level", DataType.INT64)
+
+            bm25_function = Function(
+                name="text_bm25_emb",
+                function_type=FunctionType.BM25,
+                input_field_names=["text"],
+                output_field_names=["sparse_embedding"],
+            )
+            schema.add_function(bm25_function)
+
+            index_params = client.prepare_index_params()
+            index_params.add_index(
+                field_name="dense_embedding",
+                index_type="HNSW",
+                metric_type="IP",
+                params={"M": 16, "efConstruction": 256},
+            )
+            index_params.add_index(
+                field_name="sparse_embedding",
+                index_type="SPARSE_INVERTED_INDEX",
+                metric_type="BM25",
+                params={"drop_ratio_build": 0.2},
+            )
+            try:
+                client.create_collection(
+                    collection_name=collection_name,
+                    schema=schema,
+                    index_params=index_params,
+                )
+            except Exception as e:
+                # 仅当 collection 已实际建好时才吞掉异常(典型场景:Milvus Lite on Windows
+                # 重命名 manifest.json 时偶发 WinError 183)。此时索引是否齐全会由后续
+                # _ensure_indexes 对账补建,避免历史那种"半建状态被永久隐藏"的缺陷。
+                # collection 未建成功时必须 raise,不能掩盖真正的创建失败。
+                if not client.has_collection(collection_name):
+                    raise
+
+        MilvusStore._ensure_indexes(client, collection_name)
+        MilvusStore._ensure_loaded(client, collection_name)
+
+    @staticmethod
+    def _ensure_indexes(client: MilvusClient, collection_name: str) -> None:
+        """索引对账:dense_embedding 与 sparse_embedding 任一缺失则单独补建。
+
+        默认情况下索引名等于字段名(create_index 未显式指定 index_name),
+        因此直接用字段名判断是否已存在。
+        """
         try:
-            client.create_collection(
-                collection_name=collection_name,
-                schema=schema,
-                index_params=index_params,
+            existing = set(client.list_indexes(collection_name))
+        except Exception:
+            existing = set()
+
+        if "dense_embedding" not in existing:
+            dense_params = client.prepare_index_params()
+            dense_params.add_index(
+                field_name="dense_embedding",
+                index_type="HNSW",
+                metric_type="IP",
+                params={"M": 16, "efConstruction": 256},
+            )
+            client.create_index(collection_name=collection_name, index_params=dense_params)
+
+        if "sparse_embedding" not in existing:
+            sparse_params = client.prepare_index_params()
+            sparse_params.add_index(
+                field_name="sparse_embedding",
+                index_type="SPARSE_INVERTED_INDEX",
+                metric_type="BM25",
+                params={"drop_ratio_build": 0.2},
             )
-        except Exception as e:
-            # Milvus Lite on Windows 在 create_index 后重命名 manifest.json 时偶发
-            # WinError 183,但集合与索引实际已创建成功,因此若集合已存在则忽略。
-            if client.has_collection(collection_name):
-                return
-            raise
+            client.create_index(collection_name=collection_name, index_params=sparse_params)
 
     def init_collection(self, dense_dim: int | None = None) -> None:
         if dense_dim is None:
@@ -183,7 +240,7 @@ class MilvusStore:
                 offset=offset,
             )
 
-        return self._run(_query)
+        return self._run_read(_query)
 
     def query_all(self, filter_expr: str = "", output_fields: list[str] | None = None) -> list:
         """分页拉取;单次 session 内完成,避免每页新建连接。"""
@@ -209,7 +266,7 @@ class MilvusStore:
                 offset += len(batch)
             return out
 
-        return self._run(_query_all)
+        return self._run_read(_query_all)
 
     def get_chunks_by_ids(self, chunk_ids: list[str]) -> list[dict]:
         ids = [item for item in chunk_ids if item]
@@ -278,7 +335,7 @@ class MilvusStore:
                 output_fields=output_fields,
             )
 
-        results = self._run(_search)
+        results = self._run_read(_search)
         formatted_results = []
         for hits in results:
             for hit in hits:
@@ -329,7 +386,7 @@ class MilvusStore:
                 filter=filter_expr,
             )
 
-        results = self._run(_search)
+        results = self._run_read(_search)
         formatted_results = []
         for hits in results:
             for hit in hits:

+ 8 - 1
backend/rag-ai-bridge/server.py

@@ -328,6 +328,10 @@ def text2cypher(req: Text2CypherRequest):
             "Generate one read-only Cypher query for the question. "
             "Return only Cypher, without explanation. Never use CREATE, MERGE, DELETE, SET, REMOVE, DROP, LOAD CSV or CALL.\n"
             "Return nodes, relationships, or paths that carry the requested properties; do not return only scalar projections.\n"
+            "Every node pattern MUST declare exactly one label copied EXACTLY from allowedLabels (preserve uppercase, underscore, spelling). "
+            "NEVER invent, translate, pluralize, or guess labels from the question text — every label in the Cypher MUST appear verbatim in allowedLabels. "
+            "If a concept has no exact match, pick the closest allowedLabels entry and reuse its exact spelling. "
+            "Repeating a previously-bound variable is fine, but any new node variable must carry a label.\n"
             f"Schema:\n{req.schema}\n"
             f"Allowed labels: {req.allowedLabels}\n"
             f"Allowed relationships: {req.allowedRelationships}\n"
@@ -420,10 +424,13 @@ def repair(req: RepairRequest):
     language = req.language.upper()
     prompt = f"""Repair this read-only {language} query after EXPLAIN failed. Return only the corrected query.
 Use only the supplied schema. Make exactly one statement. Do not use write operations.
+The failed query used a label or relationship type that is NOT in the Schema below — open the Schema, find the closest matching entry, and copy its name EXACTLY (preserve uppercase, underscore, spelling).
+NEVER invent, translate, pluralize, or guess labels from the question text; every label/relationship MUST appear verbatim in the Schema.
+Every node pattern MUST declare exactly one label from the Schema — write (p:Person), never (p) or ().
 Question: {req.question}\nSchema: {req.schemaText}\nFailed query: {req.query}\nError: {_short_error(req.error)}
 Maximum rows: {req.maxRows}; maximum graph depth: {req.maxDepth}."""
     try:
-        response = _openai_client(20).chat.completions.create(model=SETTINGS["llm"]["model"],
+        response = _openai_client(SETTINGS["llm"]["timeout"]).chat.completions.create(model=SETTINGS["llm"]["model"],
             messages=[{"role":"user","content":prompt}], temperature=SETTINGS["llm"]["temperature"],
             max_tokens=SETTINGS["llm"]["generation_max_tokens"])
         fixed = _read_only(_extract_code(response.choices[0].message.content or "", language), language)

+ 39 - 2
backend/src/main/java/com/agent/management/config/RagAiBridgeProperties.java

@@ -1,7 +1,44 @@
 package com.agent.management.config;
+
 import lombok.Data;
 import org.springframework.boot.context.properties.ConfigurationProperties;
 import org.springframework.stereotype.Component;
+
 import java.time.Duration;
-@Data @Component @ConfigurationProperties(prefix="rag.ai-bridge")
-public class RagAiBridgeProperties { private boolean enabled=false; private String baseUrl="http://127.0.0.1:18733"; private Duration timeout=Duration.ofSeconds(300); }
+
+/**
+ * RAG AI Bridge 配置项(rag.ai-bridge.*)
+ *
+ * 用于自然语言转 SQL/Cypher 的 Python Bridge 服务。
+ * LLM 配置(api-key/base-url/model)不在本类中,统一从 spring.ai.openai.* 读取,
+ * 由 RagAiBridgeProcessManager 通过环境变量注入子进程,实现 application.yml 单一配置源。
+ */
+@Data
+@Component
+@ConfigurationProperties(prefix = "rag.ai-bridge")
+public class RagAiBridgeProperties {
+
+    /** 是否启用 RAG AI Bridge */
+    private boolean enabled = false;
+
+    /** HTTP 调用超时(RagAiBridgeClient 用) */
+    private Duration timeout = Duration.ofSeconds(300);
+
+    /** Bridge 监听主机 */
+    private String host = "127.0.0.1";
+
+    /** Bridge 监听端口 */
+    private int port = 18733;
+
+    /** Python 可执行文件路径 */
+    private String pythonPath = "python";
+
+    /** Bridge 入口脚本路径(相对 user.dir,与 embedding-bridge / hermes-bridge 约定一致) */
+    private String scriptPath = "rag-ai-bridge/server.py";
+
+    /** 启动超时(秒) */
+    private int startupTimeout = 60;
+
+    /** 健康检查间隔(秒),0 表示不检查 */
+    private int healthCheckInterval = 60;
+}

+ 83 - 0
backend/src/main/java/com/agent/management/engine/ContextPromptHelper.java

@@ -0,0 +1,83 @@
+package com.agent.management.engine;
+
+import java.util.Map;
+
+/**
+ * 工作流上下文变量 → 系统提示词的工具类。
+ *
+ * <p>用于把前置节点输出的 variables 自动注入到 LLM/智能操作/技能等节点的系统提示词中,
+ * 让模型在用户未显式声明 {{var}} 占位符时也能"看到"前置节点的输出。</p>
+ *
+ * <p>DRY 集中点:5 个执行器(LlmExecutor / SmartActionExecutor / AgentExecutor /
+ * HermesSmartActionExecutor / HermesAgentExecutor)共享同一份格式化逻辑。</p>
+ */
+public final class ContextPromptHelper {
+
+    /** 单个变量值的截断上限,避免长内容(如整篇文档)撑爆 prompt */
+    private static final int MAX_VALUE_LENGTH = 2000;
+
+    private ContextPromptHelper() {
+    }
+
+    /**
+     * 把工作流上下文的 variables 格式化为系统提示词片段。
+     *
+     * @param context 工作流上下文(可能为 null,调用方方便起见)
+     * @return 格式化后的提示词;如果上下文为空,返回空串
+     */
+    public static String buildContextSystemPrompt(WorkflowContext context) {
+        if (context == null) {
+            return "";
+        }
+        Map<String, Object> vars = context.getVariables();
+        if (vars == null || vars.isEmpty()) {
+            return "";
+        }
+        StringBuilder sb = new StringBuilder();
+        sb.append("=== 工作流上下文(来自前置节点的输出变量) ===\n");
+        sb.append("以下是当前可用的变量值,可在回答中引用:\n\n");
+        for (Map.Entry<String, Object> e : vars.entrySet()) {
+            Object v = e.getValue();
+            if (v == null) {
+                continue;
+            }
+            sb.append("- ").append(e.getKey()).append(":").append(formatValue(v)).append("\n");
+        }
+        return sb.toString();
+    }
+
+    /**
+     * 把上下文系统提示词合并到用户自定义的系统提示词后。
+     * 两者皆空时返回空串;其中一方为空时返回另一方。
+     *
+     * @param userSystemPrompt 节点本身声明的 systemPrompt(已渲染)
+     * @param context          工作流上下文
+     * @return 合并后的完整 systemPrompt
+     */
+    public static String merge(String userSystemPrompt, WorkflowContext context) {
+        String contextPrompt = buildContextSystemPrompt(context);
+        boolean userEmpty = userSystemPrompt == null || userSystemPrompt.isBlank();
+        boolean ctxEmpty = contextPrompt.isEmpty();
+        if (userEmpty && ctxEmpty) {
+            return "";
+        }
+        if (userEmpty) {
+            return contextPrompt;
+        }
+        if (ctxEmpty) {
+            return userSystemPrompt;
+        }
+        return userSystemPrompt + "\n\n" + contextPrompt;
+    }
+
+    private static String formatValue(Object value) {
+        if (value == null) {
+            return "";
+        }
+        String s = value.toString();
+        if (s.length() <= MAX_VALUE_LENGTH) {
+            return s;
+        }
+        return s.substring(0, MAX_VALUE_LENGTH) + "...(已截断,共 " + s.length() + " 字符)";
+    }
+}

+ 2 - 2
backend/src/main/java/com/agent/management/engine/WorkflowEngine.java

@@ -44,8 +44,8 @@ public class WorkflowEngine {
             }
     );
 
-    /** 工作流最大执行时间(5 分钟) */
-    private static final long MAX_EXECUTION_SECONDS = 300;
+    /** 工作流最大执行时间:无限制(Long.MAX_VALUE 秒 ≈ 2920 亿年,scheduler.schedule 实际永不触发) */
+    private static final long MAX_EXECUTION_SECONDS = Long.MAX_VALUE;
 
     private final ScheduledExecutorService scheduler = Executors.newSingleThreadScheduledExecutor(r -> {
         Thread t = new Thread(r, "workflow-timeout");

+ 126 - 46
backend/src/main/java/com/agent/management/engine/WorkflowLevelExecutor.java

@@ -5,6 +5,7 @@ import com.agent.management.model.entity.WorkflowRunNode;
 import com.fasterxml.jackson.databind.JsonNode;
 import com.fasterxml.jackson.databind.ObjectMapper;
 import com.fasterxml.jackson.datatype.jsr310.JavaTimeModule;
+import jakarta.annotation.PreDestroy;
 import lombok.extern.slf4j.Slf4j;
 import org.springframework.stereotype.Component;
 import org.springframework.web.servlet.mvc.method.annotation.SseEmitter;
@@ -16,6 +17,10 @@ import java.util.LinkedHashMap;
 import java.util.List;
 import java.util.Map;
 import java.util.Set;
+import java.util.concurrent.CompletableFuture;
+import java.util.concurrent.ExecutorService;
+import java.util.concurrent.Executors;
+import java.util.concurrent.TimeUnit;
 
 /**
  * 按 DAG 拓扑层级逐层执行节点。
@@ -32,6 +37,18 @@ public class WorkflowLevelExecutor {
 
     private final Map<String, NodeExecutor> executorMap;
     private final SseEventBus eventBus;
+    /**
+     * 同层节点并行执行的线程池。独立于 WorkflowEngine.executor(工作流主线程池),
+     * 避免层内并行与层间串行争用同一池导致死锁。
+     */
+    private final ExecutorService nodeExecutor = Executors.newFixedThreadPool(
+            Math.max(8, Runtime.getRuntime().availableProcessors() * 2),
+            r -> {
+                Thread t = new Thread(r, "workflow-node-executor");
+                t.setDaemon(true);
+                return t;
+            }
+    );
 
     public WorkflowLevelExecutor(Map<String, NodeExecutor> injected, SseEventBus eventBus) {
         // Spring 注入 Map<String, NodeExecutor> 时 key 是 bean name(如 "userInputExecutor"),
@@ -46,6 +63,19 @@ public class WorkflowLevelExecutor {
         log.info("[WorkflowLevelExecutor] 已注册节点执行器: {}", this.executorMap.keySet());
     }
 
+    @PreDestroy
+    public void shutdown() {
+        nodeExecutor.shutdown();
+        try {
+            if (!nodeExecutor.awaitTermination(5, TimeUnit.SECONDS)) {
+                nodeExecutor.shutdownNow();
+            }
+        } catch (InterruptedException e) {
+            Thread.currentThread().interrupt();
+            nodeExecutor.shutdownNow();
+        }
+    }
+
     /**
      * 一次工作流执行的产物:最终输出 + 每节点记录 + 是否失败 + 错误信息
      */
@@ -88,58 +118,50 @@ public class WorkflowLevelExecutor {
         Map<String, Object> finalOutputs = new LinkedHashMap<>();
         List<WorkflowRunNode> nodeRecords = new ArrayList<>();
         int sortOrder = 0;
-        boolean failed = false;
-        String errorMsg = null;
 
         for (int levelIdx = 0; levelIdx < dag.getLevels().size(); levelIdx++) {
             List<String> level = dag.getLevels().get(levelIdx);
             log.debug("[WorkflowEngine] 执行第 {} 层, {} 个节点, 活跃: {}", levelIdx, level.size(),
                     level.stream().filter(activeNodes::contains).count());
 
+            // 1. 分离活跃与非活跃节点;非活跃节点直接标记为 SKIPPED(保持原 sortOrder 顺序)
+            List<String> activeInLevel = new ArrayList<>();
             for (String nodeId : level) {
-                String nodeType = dag.getNodeTypeMap().get(nodeId);
-                JsonNode nodeData = dag.getNodeDataMap().get(nodeId);
-                String label = nodeData.path("label").asText(nodeId);
-
-                if (!activeNodes.contains(nodeId)) {
+                if (activeNodes.contains(nodeId)) {
+                    activeInLevel.add(nodeId);
+                } else {
+                    String nodeType = dag.getNodeTypeMap().get(nodeId);
+                    JsonNode nodeData = dag.getNodeDataMap().get(nodeId);
+                    String label = nodeData.path("label").asText(nodeId);
                     safeSend(emitter, WorkflowRunEvent.nodeResult(runId, NodeExecutionResult.skipped(nodeId)));
-                    nodeRecords.add(buildNodeRecord(runRecordId, nodeId, nodeType, label, "SKIPPED", null, null, null, null, null, sortOrder++));
-                    continue;
-                }
-
-                NodeExecutor nodeExecutor = executorMap.get(nodeType);
-                if (nodeExecutor == null) {
-                    String errMsg = "未知节点类型: " + nodeType;
-                    safeSend(emitter, WorkflowRunEvent.nodeResult(runId, NodeExecutionResult.failed(nodeId, errMsg)));
-                    safeSend(emitter, WorkflowRunEvent.workflowError(runId, errMsg));
-                    nodeRecords.add(buildNodeRecord(runRecordId, nodeId, nodeType, label, "FAILED", null, errMsg, null, null, null, sortOrder++));
-                    return new ExecutionOutcome(finalOutputs, nodeRecords, true, errMsg);
+                    nodeRecords.add(buildNodeRecord(runRecordId, nodeId, nodeType, label,
+                            "SKIPPED", null, null, null, null, null, sortOrder++));
                 }
+            }
 
-                safeSend(emitter, WorkflowRunEvent.nodeRunning(runId, nodeId));
-
-                // 前置条件校验
-                String precheckError = checkPreconditions(nodeData, context);
-                if (precheckError != null) {
-                    log.warn("[WorkflowEngine] 节点 {} 前置条件不满足: {}", nodeId, precheckError);
-                    safeSend(emitter, WorkflowRunEvent.nodeResult(runId, NodeExecutionResult.failed(nodeId, precheckError)));
-                    nodeRecords.add(buildNodeRecord(runRecordId, nodeId, nodeType, label, "FAILED", null, precheckError, null, null, null, sortOrder++));
+            if (activeInLevel.isEmpty()) continue;
 
-                    String failStrategy = nodeData.path("failStrategy").asText("abort");
-                    if ("skip".equals(failStrategy)) {
-                        continue;
-                    }
-                    String abortMsg = "节点 " + nodeId + " 前置条件不满足: " + precheckError;
-                    safeSend(emitter, WorkflowRunEvent.workflowError(runId, abortMsg));
-                    return new ExecutionOutcome(finalOutputs, nodeRecords, true, abortMsg);
-                }
+            // 2. 同层活跃节点并行执行:每个节点提交到 nodeExecutor 线程池
+            List<CompletableFuture<NodeExecutionResult>> futures = new ArrayList<>(activeInLevel.size());
+            for (String nodeId : activeInLevel) {
+                final String finalNodeId = nodeId;
+                futures.add(CompletableFuture.supplyAsync(
+                        () -> executeOneNode(finalNodeId, dag, context, runId, emitter),
+                        nodeExecutor));
+            }
+            // 等待当前层所有并行节点完成(任一节点抛出的异常都被封装为 failed 结果,不会从 join 传播)
+            CompletableFuture.allOf(futures.toArray(new CompletableFuture[0])).join();
+
+            // 3. 串行处理结果:按提交顺序写 nodeRecords / sortOrder,按 failStrategy 决定是否 abort
+            //    结果处理阶段是单线程,无需同步 nodeRecords / finalOutputs / sortOrder
+            String abortMsg = null;
+            for (int i = 0; i < activeInLevel.size(); i++) {
+                String nodeId = activeInLevel.get(i);
+                NodeExecutionResult result = joinResult(futures.get(i), nodeId);
+                String nodeType = dag.getNodeTypeMap().get(nodeId);
+                JsonNode nodeData = dag.getNodeDataMap().get(nodeId);
+                String label = nodeData.path("label").asText(nodeId);
 
-                NodeExecutionResult result = nodeExecutor.execute(nodeId, nodeData, context);
-                if (result.getOutput() != null) {
-                    context.setNodeOutput(nodeId, result.getOutput());
-                }
-                // 节点完成后采集完整 variables 快照(debug 视图与运行历史回放)
-                result = result.withContextSnapshot(context.snapshotVariables());
                 safeSend(emitter, WorkflowRunEvent.nodeResult(runId, result));
                 nodeRecords.add(buildNodeRecord(runRecordId, nodeId, nodeType, label,
                         result.getStatus().name(), result.getOutput(), result.getError(),
@@ -149,10 +171,9 @@ public class WorkflowLevelExecutor {
                     String failStrategy = nodeData.path("failStrategy").asText("abort");
                     if ("skip".equals(failStrategy)) {
                         log.info("[WorkflowEngine] 节点 {} 执行失败,策略为跳过,继续工作流", nodeId);
-                    } else {
-                        String failMsg = "节点 " + nodeId + " 执行失败: " + result.getError();
-                        safeSend(emitter, WorkflowRunEvent.workflowError(runId, failMsg));
-                        return new ExecutionOutcome(finalOutputs, nodeRecords, true, failMsg);
+                    } else if (abortMsg == null) {
+                        // 同层多个失败时,仅记录第一个 abort 原因;继续处理剩余结果以保证 nodeRecords 完整
+                        abortMsg = "节点 " + nodeId + " 执行失败: " + result.getError();
                     }
                 }
 
@@ -163,9 +184,63 @@ public class WorkflowLevelExecutor {
                     finalOutputs.putAll(result.getOutput());
                 }
             }
+
+            if (abortMsg != null) {
+                safeSend(emitter, WorkflowRunEvent.workflowError(runId, abortMsg));
+                return new ExecutionOutcome(finalOutputs, nodeRecords, true, abortMsg);
+            }
+        }
+
+        return new ExecutionOutcome(finalOutputs, nodeRecords, false, null);
+    }
+
+    /**
+     * 执行单个节点(线程池 worker 中调用)。
+     * 把节点类型查找、前置条件校验、executor.execute、上下文写入等放在同一个并行任务里。
+     * 异常一律封装为 failed 结果返回,不向上抛(避免中断 CompletableFuture.allOf)。
+     */
+    private NodeExecutionResult executeOneNode(String nodeId, DagResolver.ResolvedDag dag,
+                                               WorkflowContext context, String runId, SseEmitter emitter) {
+        String nodeType = dag.getNodeTypeMap().get(nodeId);
+        JsonNode nodeData = dag.getNodeDataMap().get(nodeId);
+
+        NodeExecutor executor = executorMap.get(nodeType);
+        if (executor == null) {
+            String errMsg = "未知节点类型: " + nodeType;
+            log.warn("[WorkflowEngine] 节点 {} {}", nodeId, errMsg);
+            return NodeExecutionResult.failed(nodeId, errMsg);
+        }
+
+        safeSend(emitter, WorkflowRunEvent.nodeRunning(runId, nodeId));
+
+        String precheckError = checkPreconditions(nodeData, context);
+        if (precheckError != null) {
+            log.warn("[WorkflowEngine] 节点 {} 前置条件不满足: {}", nodeId, precheckError);
+            return NodeExecutionResult.failed(nodeId, precheckError);
+        }
+
+        try {
+            NodeExecutionResult result = executor.execute(nodeId, nodeData, context);
+            if (result.getOutput() != null) {
+                context.setNodeOutput(nodeId, result.getOutput());
+            }
+            // 节点完成后采集完整 variables 快照(debug 视图与运行历史回放)
+            return result.withContextSnapshot(context.snapshotVariables());
+        } catch (Exception e) {
+            log.error("[WorkflowEngine] 节点 {} 执行抛出异常: {}", nodeId, e.getMessage(), e);
+            return NodeExecutionResult.failed(nodeId, "节点执行异常: " + e.getMessage());
         }
+    }
 
-        return new ExecutionOutcome(finalOutputs, nodeRecords, failed, errorMsg);
+    /**
+     * 从 CompletableFuture 取出结果;future 内部异常一律降级为 failed。
+     */
+    private static NodeExecutionResult joinResult(CompletableFuture<NodeExecutionResult> future, String nodeId) {
+        try {
+            return future.join();
+        } catch (Exception e) {
+            return NodeExecutionResult.failed(nodeId, "节点并行执行异常: " + e.getMessage());
+        }
     }
 
     /**
@@ -259,7 +334,12 @@ public class WorkflowLevelExecutor {
         return missing.isEmpty() ? null : "前置条件不满足: " + String.join("; ", missing);
     }
 
-    private void safeSend(SseEmitter emitter, WorkflowRunEvent event) {
+    /**
+     * 安全发送 SSE 事件
+     * <p>加 synchronized:同层多节点并行执行时,多个 worker 线程会并发调用 safeSend
+     * (nodeRunning / nodeStream / nodeResult),而 SseEmitter.send 非线程安全。</p>
+     */
+    private synchronized void safeSend(SseEmitter emitter, WorkflowRunEvent event) {
         // 同时发布到事件总线(外部 API 通过 /stream 订阅消费)
         eventBus.publish(event.getRunId(), event.getType(), event);
         try {

+ 3 - 1
backend/src/main/java/com/agent/management/engine/executor/AgentExecutor.java

@@ -65,8 +65,10 @@ public class AgentExecutor implements NodeExecutor {
         log.info("[Agent] 节点 {} 执行 Skill: {}, 用户消息长度={}", nodeId, agentId, userMessage.length());
 
         try {
+            // 把前置节点 variables 追加到 SKILL.md 内容之后,作为完整系统提示
+            String fullSystem = ContextPromptHelper.merge(skillContent, context);
             var request = client.prompt()
-                    .system(skillContent)
+                    .system(fullSystem)
                     .user(userMessage);
             String result = request.call().content();
 

+ 2 - 1
backend/src/main/java/com/agent/management/engine/executor/HermesAgentExecutor.java

@@ -58,7 +58,8 @@ public class HermesAgentExecutor implements NodeExecutor {
         String userMessage;
         try {
             SkillDTO skill = skillService.getSkill(agentId);
-            systemPrompt = skillService.getSkillFullContent(agentId);
+            // 把前置节点 variables 追加到 SKILL.md 内容之后,作为完整系统提示
+            systemPrompt = ContextPromptHelper.merge(skillService.getSkillFullContent(agentId), context);
             userMessage = buildUserMessage(skill, context);
         } catch (Exception e) {
             return NodeExecutionResult.failed(nodeId, "加载 Skill 失败: " + e.getMessage());

+ 4 - 1
backend/src/main/java/com/agent/management/engine/executor/HermesSmartActionExecutor.java

@@ -59,7 +59,10 @@ public class HermesSmartActionExecutor implements NodeExecutor {
             String workingDir = context.getWorkingDir() != null ? context.getWorkingDir().toString() : null;
             String sessionId = context.getRunId() + "_" + nodeId;
             Map<String, String> modelConfig = resolveModelConfig(data);
-            HermesRunResult runResult = bridgeClient.run(null, actionPrompt, maxIterations,
+            // 把前置节点 variables 作为系统提示注入 Hermes Bridge(原 null 不再使用)
+            String contextSystemPrompt = ContextPromptHelper.buildContextSystemPrompt(context);
+            HermesRunResult runResult = bridgeClient.run(contextSystemPrompt.isEmpty() ? null : contextSystemPrompt,
+                    actionPrompt, maxIterations,
                     hermesHome, workingDir, nodeId, context.getStreamSink(), sessionId, modelConfig);
 
             if (runResult.getFinalText() == null || runResult.getFinalText().isBlank()) {

+ 257 - 0
backend/src/main/java/com/agent/management/engine/executor/KnowledgeRetrievalExecutor.java

@@ -0,0 +1,257 @@
+package com.agent.management.engine.executor;
+
+import com.agent.management.engine.NodeExecutionResult;
+import com.agent.management.engine.NodeExecutor;
+import com.agent.management.engine.TemplateRenderer;
+import com.agent.management.engine.WorkflowContext;
+import com.agent.management.model.entity.DataSource;
+import com.agent.management.model.entity.GraphSource;
+import com.agent.management.rag.document.DocumentRagRetriever;
+import com.agent.management.rag.graph.GraphRagRetriever;
+import com.agent.management.rag.kb.KnowledgeBaseRagRetriever;
+import com.agent.management.rag.model.KnowledgeBaseRagRequest;
+import com.agent.management.rag.model.KnowledgeBaseRagResult;
+import com.agent.management.rag.model.RagEvidence;
+import com.agent.management.rag.model.RagQuery;
+import com.agent.management.rag.model.RagRetrievalResult;
+import com.agent.management.rag.model.RagSourceType;
+import com.agent.management.rag.structured.StructuredDataRagRetriever;
+import com.agent.management.service.DataSourceService;
+import com.agent.management.service.GraphSourceService;
+import com.fasterxml.jackson.databind.JsonNode;
+import lombok.RequiredArgsConstructor;
+import lombok.extern.slf4j.Slf4j;
+import org.springframework.stereotype.Component;
+
+import java.util.ArrayList;
+import java.util.LinkedHashMap;
+import java.util.List;
+import java.util.Map;
+
+/**
+ * 知识库检索节点执行器
+ *
+ * <p>支持四种检索来源:</p>
+ * <ul>
+ *   <li><b>document</b>:文档数据库(DocumentRagRetriever)。sourceIds 支持多个文档 ID,
+ *       allDocuments=true 时 sourceIds 留空,由 retriever 检索全部文档。</li>
+ *   <li><b>structured</b>:结构化数据库(StructuredDataRagRetriever)。retriever 仅取 sourceIds[0],
+ *       因此多数据源时本执行器循环调用每个 datasourceId 并合并 evidences。</li>
+ *   <li><b>graph</b>:知识图谱库(GraphRagRetriever)。retriever 仅取 sourceIds[0],
+ *       多图谱时同样循环调用合并。</li>
+ *   <li><b>hybrid</b>:混合检索(KnowledgeBaseRagRetriever)。按知识库绑定关系自动并行调用上述三类 retriever。</li>
+ * </ul>
+ *
+ * <p>query 字段支持 {@code {{变量名}}} 模板渲染,可引用上游节点输出。</p>
+ *
+ * <p>输出写入工作流上下文:</p>
+ * <ul>
+ *   <li><b>evidences</b>:{@code List<RagEvidence>}</li>
+ *   <li><b>evidenceCount</b>:int</li>
+ *   <li><b>sourceType</b>:DOCUMENT / STRUCTURED_DATA / GRAPH / HYBRID</li>
+ *   <li><b>diagnostics</b>:{@code Map<String,Object>}</li>
+ * </ul>
+ */
+@Slf4j
+@Component
+@RequiredArgsConstructor
+public class KnowledgeRetrievalExecutor implements NodeExecutor {
+
+    private final DocumentRagRetriever documentRetriever;
+    private final StructuredDataRagRetriever structuredRetriever;
+    private final GraphRagRetriever graphRetriever;
+    private final KnowledgeBaseRagRetriever knowledgeBaseRetriever;
+    private final DataSourceService dataSourceService;
+    private final GraphSourceService graphSourceService;
+
+    @Override
+    public String getType() {
+        return "knowledgeRetrieval";
+    }
+
+    @Override
+    public NodeExecutionResult execute(String nodeId, JsonNode data, WorkflowContext context) {
+        String source = data.path("source").asText("document");
+        String rawQuery = data.path("query").asText("");
+        String query = TemplateRenderer.render(rawQuery, context.getVariables()).trim();
+
+        if (query.isEmpty()) {
+            return NodeExecutionResult.failed(nodeId, "知识库检索节点的检索语句为空");
+        }
+
+        int topK = data.path("topK").asInt(5);
+        if (topK <= 0) topK = 5;
+
+        try {
+            Map<String, Object> output = switch (source) {
+                case "document" -> retrieveDocument(data, query, topK);
+                case "structured" -> retrieveStructured(data, query, topK);
+                case "graph" -> retrieveGraph(data, query, topK);
+                case "hybrid" -> retrieveHybrid(data, query, topK);
+                default -> {
+                    log.warn("[KnowledgeRetrieval] 节点 {} 未知 source={}, 降级为 document", nodeId, source);
+                    yield retrieveDocument(data, query, topK);
+                }
+            };
+
+            log.info("[KnowledgeRetrieval] 节点 {} 完成, source={}, evidences={}",
+                    nodeId, source, output.get("evidenceCount"));
+            return NodeExecutionResult.success(nodeId, output);
+        } catch (IllegalArgumentException e) {
+            return NodeExecutionResult.failed(nodeId, e.getMessage());
+        } catch (Exception e) {
+            log.error("[KnowledgeRetrieval] 节点 {} 检索失败: {}", nodeId, e.getMessage(), e);
+            return NodeExecutionResult.failed(nodeId, "知识库检索失败:" + e.getMessage());
+        }
+    }
+
+    // =========================== 检索实现 ===========================
+
+    private Map<String, Object> retrieveDocument(JsonNode data, String query, int topK) {
+        boolean allDocuments = data.path("allDocuments").asBoolean(true);
+        List<String> documentIds = stringList(data, "documentIds");
+
+        if (!allDocuments && documentIds.isEmpty()) {
+            throw new IllegalArgumentException("未选择文档范围,请勾选文档或改为「在所有文档中检索」");
+        }
+
+        RagQuery rq = buildQuery(query, topK, allDocuments ? null : documentIds);
+        rq.setFilters(Map.of("mode", "hybrid"));
+
+        RagRetrievalResult result = documentRetriever.retrieve(rq);
+        return toOutputMap(result);
+    }
+
+    private Map<String, Object> retrieveStructured(JsonNode data, String query, int topK) {
+        List<String> datasourceIds = resolveScopedSourceIds(
+                data, "allDatasources", "datasourceIds",
+                () -> dataSourceService.listAll().stream().map(DataSource::getId).map(String::valueOf).toList());
+        if (datasourceIds.isEmpty()) {
+            throw new IllegalArgumentException("未选择结构化数据源,且系统中无可用数据源");
+        }
+
+        List<RagEvidence> evidences = new ArrayList<>();
+        Map<String, Object> diagnostics = new LinkedHashMap<>();
+        for (String dsId : datasourceIds) {
+            RagQuery rq = buildQuery(query, topK, List.of(dsId));
+            Map<String, Object> filters = new LinkedHashMap<>();
+            filters.put("allowTextToSql", true);
+            filters.put("maxRows", topK);
+            rq.setFilters(filters);
+            try {
+                RagRetrievalResult result = structuredRetriever.retrieve(rq);
+                mergeRetrieval(result, "datasource:" + dsId, evidences, diagnostics);
+            } catch (Exception e) {
+                diagnostics.put("error:datasource:" + dsId, e.getMessage());
+                log.warn("[KnowledgeRetrieval] 结构化数据源 {} 检索失败: {}", dsId, e.getMessage());
+            }
+        }
+        return aggregate(RagSourceType.STRUCTURED_DATA, evidences, diagnostics);
+    }
+
+    private Map<String, Object> retrieveGraph(JsonNode data, String query, int topK) {
+        List<String> graphIds = resolveScopedSourceIds(
+                data, "allGraphSources", "graphSourceIds",
+                () -> graphSourceService.listAll().stream().map(GraphSource::getId).map(String::valueOf).toList());
+        if (graphIds.isEmpty()) {
+            throw new IllegalArgumentException("未选择图谱数据源,且系统中无可用图谱");
+        }
+
+        List<RagEvidence> evidences = new ArrayList<>();
+        Map<String, Object> diagnostics = new LinkedHashMap<>();
+        for (String gId : graphIds) {
+            RagQuery rq = buildQuery(query, topK, List.of(gId));
+            rq.setFilters(Map.of("allowTextToCypher", true));
+            try {
+                RagRetrievalResult result = graphRetriever.retrieve(rq);
+                mergeRetrieval(result, "graph:" + gId, evidences, diagnostics);
+            } catch (Exception e) {
+                diagnostics.put("error:graph:" + gId, e.getMessage());
+                log.warn("[KnowledgeRetrieval] 图谱 {} 检索失败: {}", gId, e.getMessage());
+            }
+        }
+        return aggregate(RagSourceType.GRAPH, evidences, diagnostics);
+    }
+
+    private Map<String, Object> retrieveHybrid(JsonNode data, String query, int topK) {
+        long kbId = data.path("knowledgeBaseId").asLong(0);
+        if (kbId <= 0) {
+            throw new IllegalArgumentException("混合检索需要指定有效的 knowledgeBaseId");
+        }
+
+        KnowledgeBaseRagRequest req = new KnowledgeBaseRagRequest();
+        req.setKnowledgeBaseId(kbId);
+        req.setQuery(query);
+        req.setTopK(topK);
+        KnowledgeBaseRagResult result = knowledgeBaseRetriever.retrieve(req);
+
+        Map<String, Object> output = new LinkedHashMap<>();
+        output.put("evidences", result.getEvidences());
+        output.put("evidenceCount", result.getEvidences().size());
+        output.put("sourceType", "HYBRID");
+        output.put("knowledgeBaseId", kbId);
+        output.put("enabledSources", result.getEnabledSources());
+        output.put("diagnostics", result.getDiagnostics() == null ? Map.of() : result.getDiagnostics());
+        return output;
+    }
+
+    // =========================== 工具 ===========================
+
+    private static RagQuery buildQuery(String query, int topK, List<String> sourceIds) {
+        RagQuery rq = new RagQuery();
+        rq.setQuery(query);
+        rq.setTopK(topK);
+        if (sourceIds != null && !sourceIds.isEmpty()) rq.setSourceIds(sourceIds);
+        return rq;
+    }
+
+    /**
+     * 解析 sourceIds:若 allFlag=true,则调用 allSupplier 取全部;否则取显式 IDs。
+     * 用于结构化和图谱的「在所有数据源中检索」语义。
+     */
+    private static List<String> resolveScopedSourceIds(JsonNode data, String allFlag, String idsField,
+                                                        java.util.function.Supplier<List<String>> allSupplier) {
+        boolean all = data.path(allFlag).asBoolean(false);
+        if (all) return allSupplier.get();
+        return stringList(data, idsField);
+    }
+
+    private static List<String> stringList(JsonNode data, String field) {
+        JsonNode arr = data.path(field);
+        if (!arr.isArray() || arr.isEmpty()) return List.of();
+        List<String> result = new ArrayList<>();
+        for (JsonNode item : arr) {
+            String s = item.asText("");
+            if (!s.isBlank()) result.add(s);
+        }
+        return result;
+    }
+
+    private static void mergeRetrieval(RagRetrievalResult result, String diagnosticKey,
+                                        List<RagEvidence> evidences, Map<String, Object> diagnostics) {
+        if (result == null) return;
+        if (result.getEvidences() != null) evidences.addAll(result.getEvidences());
+        if (result.getDiagnostics() != null && !result.getDiagnostics().isEmpty()) {
+            diagnostics.put(diagnosticKey, result.getDiagnostics());
+        }
+    }
+
+    private static Map<String, Object> toOutputMap(RagRetrievalResult result) {
+        Map<String, Object> out = new LinkedHashMap<>();
+        out.put("evidences", result.getEvidences() == null ? List.of() : result.getEvidences());
+        out.put("evidenceCount", result.getEvidences() == null ? 0 : result.getEvidences().size());
+        out.put("sourceType", result.getSourceType() == null ? "" : result.getSourceType().name());
+        out.put("diagnostics", result.getDiagnostics() == null ? Map.of() : result.getDiagnostics());
+        return out;
+    }
+
+    private static Map<String, Object> aggregate(RagSourceType type, List<RagEvidence> evidences,
+                                                   Map<String, Object> diagnostics) {
+        Map<String, Object> out = new LinkedHashMap<>();
+        out.put("evidences", evidences);
+        out.put("evidenceCount", evidences.size());
+        out.put("sourceType", type.name());
+        out.put("diagnostics", diagnostics);
+        return out;
+    }
+}

+ 4 - 1
backend/src/main/java/com/agent/management/engine/executor/LlmExecutor.java

@@ -33,7 +33,10 @@ public class LlmExecutor implements NodeExecutor {
 
     @Override
     public NodeExecutionResult execute(String nodeId, JsonNode data, WorkflowContext context) {
-        String systemPrompt = TemplateRenderer.render(data.path("systemPrompt").asText(""), context.getVariables());
+        // 自动把前置节点输出的 variables 作为系统提示词注入,无需用户显式写 {{var}}
+        String systemPrompt = ContextPromptHelper.merge(
+                TemplateRenderer.render(data.path("systemPrompt").asText(""), context.getVariables()),
+                context);
         String userPrompt = TemplateRenderer.render(data.path("userPrompt").asText(""), context.getVariables());
 
         if (userPrompt.isEmpty()) {

+ 7 - 4
backend/src/main/java/com/agent/management/engine/executor/SmartActionExecutor.java

@@ -48,10 +48,13 @@ public class SmartActionExecutor implements NodeExecutor {
         log.info("[SmartAction] 节点 {} 开始执行, actionPrompt长度={}", nodeId, actionPrompt.length());
 
         try {
-            String result = client.prompt()
-                    .user(actionPrompt)
-                    .call()
-                    .content();
+            // 智能操作节点没有用户可配置的 systemPrompt,直接用前置节点 variables 作为系统提示
+            String contextSystemPrompt = ContextPromptHelper.buildContextSystemPrompt(context);
+            var request = client.prompt().user(actionPrompt);
+            if (!contextSystemPrompt.isEmpty()) {
+                request = request.system(contextSystemPrompt);
+            }
+            String result = request.call().content();
 
             if (result == null || result.isBlank()) {
                 return NodeExecutionResult.failed(nodeId, "智能操作返回空结果");

+ 1 - 1
backend/src/main/java/com/agent/management/rag/bridge/RagAiBridgeClient.java

@@ -28,5 +28,5 @@ public class RagAiBridgeClient {
         var response=rest.postForEntity(url(path),body,Map.class);
         return response.getBody()==null?Map.of():response.getBody();
     }
-    private String url(String path){return properties.getBaseUrl().replaceAll("/+$","")+path;}
+    private String url(String path){return "http://"+properties.getHost()+":"+properties.getPort()+path;}
 }

+ 211 - 0
backend/src/main/java/com/agent/management/rag/bridge/RagAiBridgeProcessManager.java

@@ -0,0 +1,211 @@
+package com.agent.management.rag.bridge;
+
+import com.agent.management.config.RagAiBridgeProperties;
+import jakarta.annotation.PostConstruct;
+import jakarta.annotation.PreDestroy;
+import lombok.extern.slf4j.Slf4j;
+import org.springframework.beans.factory.annotation.Value;
+import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty;
+import org.springframework.stereotype.Component;
+
+import java.io.BufferedReader;
+import java.io.File;
+import java.io.IOException;
+import java.io.InputStreamReader;
+import java.nio.charset.StandardCharsets;
+import java.util.Map;
+import java.util.concurrent.Executors;
+import java.util.concurrent.ScheduledExecutorService;
+import java.util.concurrent.TimeUnit;
+import java.util.concurrent.atomic.AtomicBoolean;
+
+/**
+ * 管理 RAG AI Bridge Python 子进程的生命周期。
+ * 仅在 rag.ai-bridge.enabled=true 时激活。
+ *
+ * LLM 配置从 spring.ai.openai.* 读取,通过环境变量注入子进程
+ * (OPENAI_API_KEY / OPENAI_BASE_URL / OPENAI_MODEL / LLM_TEMPERATURE),
+ * 实现 application.yml 作为大模型配置的单一来源。
+ */
+@Slf4j
+@Component
+@ConditionalOnProperty(name = "rag.ai-bridge.enabled", havingValue = "true")
+public class RagAiBridgeProcessManager {
+
+    private final RagAiBridgeProperties properties;
+    private final RagAiBridgeClient client;
+    /** spring.ai.openai.* 中的 LLM 配置,通过环境变量传给 Bridge 子进程 */
+    private final String llmBaseUrl;
+    private final String llmApiKey;
+    private final String llmModel;
+    private final String llmTemperature;
+
+    private Process process;
+    private final AtomicBoolean starting = new AtomicBoolean(false);
+    private ScheduledExecutorService healthScheduler;
+    /** 显式关闭标志,用于区分关闭流程与运行期异常 */
+    private volatile boolean stopped = false;
+
+    public RagAiBridgeProcessManager(
+            RagAiBridgeProperties properties,
+            RagAiBridgeClient client,
+            @Value("${spring.ai.openai.base-url:}") String llmBaseUrl,
+            @Value("${spring.ai.openai.api-key:}") String llmApiKey,
+            @Value("${spring.ai.openai.chat.options.model:}") String llmModel,
+            @Value("${spring.ai.openai.chat.options.temperature:0.3}") String llmTemperature) {
+        this.properties = properties;
+        this.client = client;
+        this.llmBaseUrl = llmBaseUrl;
+        this.llmApiKey = llmApiKey;
+        this.llmModel = llmModel;
+        this.llmTemperature = llmTemperature;
+    }
+
+    @PostConstruct
+    public void start() {
+        log.info("[RagAiBridge] 启动子进程...");
+        starting.set(true);
+
+        try {
+            String scriptPath = resolveScriptPath();
+            String pythonPath = properties.getPythonPath();
+            int port = properties.getPort();
+
+            // server.py 不支持 --port/--host 命令行参数,全部通过环境变量传
+            ProcessBuilder pb = new ProcessBuilder(pythonPath, scriptPath);
+            pb.redirectErrorStream(true);
+
+            Map<String, String> env = pb.environment();
+            env.put("PYTHONUNBUFFERED", "1");
+            // 强制 Python 子进程 stdin/stdout/stderr 使用 UTF-8(Windows 默认 cp936 会导致中文乱码)
+            env.put("PYTHONIOENCODING", "utf-8");
+            env.put("PYTHONUTF8", "1");
+            env.put("RAG_AI_BRIDGE_HOST", properties.getHost());
+            env.put("RAG_AI_BRIDGE_PORT", String.valueOf(port));
+
+            // LLM 配置注入:从 spring.ai.openai.* 读取,实现 application.yml 单一配置源
+            if (llmBaseUrl != null && !llmBaseUrl.isBlank()) {
+                env.put("OPENAI_BASE_URL", llmBaseUrl);
+            }
+            if (llmApiKey != null && !llmApiKey.isBlank()) {
+                env.put("OPENAI_API_KEY", llmApiKey);
+            }
+            if (llmModel != null && !llmModel.isBlank()) {
+                env.put("OPENAI_MODEL", llmModel);
+            }
+            if (llmTemperature != null && !llmTemperature.isBlank()) {
+                env.put("LLM_TEMPERATURE", llmTemperature);
+            }
+
+            // 日志中不输出 base_url / model / api_key 详情,避免敏感信息泄露
+            log.info("[RagAiBridge] 环境变量已注入: BASE_URL={}, MODEL={}, API_KEY={}, TEMPERATURE={}",
+                    (llmBaseUrl != null && !llmBaseUrl.isBlank() ? "已设置" : "未设置"),
+                    (llmModel != null && !llmModel.isBlank() ? "已设置" : "未设置"),
+                    (llmApiKey != null && !llmApiKey.isBlank() ? "已设置" : "未设置"),
+                    llmTemperature);
+
+            // 设置工作目录为项目根目录
+            File projectRoot = new File(System.getProperty("user.dir"));
+            pb.directory(projectRoot);
+
+            process = pb.start();
+
+            // 异步读取子进程输出(避免缓冲区满导致阻塞)
+            Thread outputReader = new Thread(() -> {
+                try (BufferedReader reader = new BufferedReader(
+                        new InputStreamReader(process.getInputStream(), StandardCharsets.UTF_8))) {
+                    String line;
+                    while ((line = reader.readLine()) != null) {
+                        log.debug("[RagAiBridge] {}", line);
+                    }
+                } catch (IOException e) {
+                    if (!stopped) {
+                        log.warn("[RagAiBridge] 输出流关闭: {}", e.getMessage());
+                    }
+                }
+            }, "rag-ai-bridge-output");
+            outputReader.setDaemon(true);
+            outputReader.start();
+
+            // 等待健康检查通过
+            long deadline = System.currentTimeMillis() + properties.getStartupTimeout() * 1000L;
+            while (System.currentTimeMillis() < deadline) {
+                Thread.sleep(1000);
+                if (client.health()) {
+                    log.info("[RagAiBridge] 子进程启动成功,端口: {}", port);
+                    starting.set(false);
+                    startHealthCheck();
+                    return;
+                }
+            }
+
+            log.error("[RagAiBridge] 子进程启动超时({}秒)", properties.getStartupTimeout());
+            starting.set(false);
+
+        } catch (Exception e) {
+            log.error("[RagAiBridge] 子进程启动失败: {}", e.getMessage(), e);
+            starting.set(false);
+        }
+    }
+
+    @PreDestroy
+    public void stop() {
+        stopped = true;
+        if (healthScheduler != null) {
+            healthScheduler.shutdownNow();
+        }
+        if (process != null && process.isAlive()) {
+            log.info("[RagAiBridge] 停止子进程...");
+            process.destroy();
+            try {
+                if (!process.waitFor(5, TimeUnit.SECONDS)) {
+                    process.destroyForcibly();
+                }
+            } catch (InterruptedException e) {
+                Thread.currentThread().interrupt();
+                process.destroyForcibly();
+            }
+            log.info("[RagAiBridge] 子进程已停止");
+        }
+    }
+
+    /**
+     * 检查 Bridge 是否就绪
+     */
+    public boolean isReady() {
+        if (starting.get()) return false;
+        if (process == null || !process.isAlive()) return false;
+        return client.health();
+    }
+
+    private void startHealthCheck() {
+        int interval = properties.getHealthCheckInterval();
+        if (interval <= 0) return;
+
+        healthScheduler = Executors.newSingleThreadScheduledExecutor(r -> {
+            Thread t = new Thread(r, "rag-ai-bridge-health-check");
+            t.setDaemon(true);
+            return t;
+        });
+
+        healthScheduler.scheduleAtFixedRate(() -> {
+            if (stopped) return;
+            if (process != null && !process.isAlive()) {
+                log.warn("[RagAiBridge] 子进程已退出,尝试重启...");
+                start();
+            }
+        }, interval, interval, TimeUnit.SECONDS);
+    }
+
+    private String resolveScriptPath() {
+        String path = properties.getScriptPath();
+        File file = new File(path);
+        if (!file.isAbsolute()) {
+            file = new File(System.getProperty("user.dir"), path);
+        }
+        if (!file.exists()) {
+            throw new IllegalStateException("RAG AI Bridge 脚本不存在: " + file.getAbsolutePath());
+        }
+        return file.getAbsolutePath();
+    }
+}

+ 2 - 2
backend/src/main/java/com/agent/management/rag/controller/KnowledgeBaseRagDebugController.java

@@ -1,2 +1,2 @@
-package com.agent.management.rag.controller; import com.agent.management.common.Result; import com.agent.management.model.entity.RagKnowledgeSourceBinding; import com.agent.management.rag.kb.KnowledgeBaseRagRetriever; import com.agent.management.rag.kb.RagKnowledgeBaseConfigService; import com.agent.management.rag.model.*; import lombok.RequiredArgsConstructor; import org.springframework.web.bind.annotation.*; import java.util.*;
-@RestController @RequestMapping("/api/rag/kb") @RequiredArgsConstructor public class KnowledgeBaseRagDebugController {private final KnowledgeBaseRagRetriever retriever;private final RagKnowledgeBaseConfigService configs;@PostMapping("/retrieve") public Result<KnowledgeBaseRagResult> retrieve(@RequestBody KnowledgeBaseRagRequest q){return Result.success(retriever.retrieve(q));}@GetMapping("/{id}/bindings") public Result<List<RagKnowledgeSourceBinding>> bindings(@PathVariable Long id){return Result.success(configs.listBindings(id));}@PutMapping("/{id}/bindings/{bindingId}/config") public Result<RagKnowledgeSourceBinding> updateConfig(@PathVariable Long id,@PathVariable Long bindingId,@RequestBody Map<String,Object> config){return Result.success(configs.updateConfig(id,bindingId,config));}}
+package com.agent.management.rag.controller; import com.agent.management.common.Result; import com.agent.management.model.entity.RagKnowledgeBase; import com.agent.management.model.entity.RagKnowledgeSourceBinding; import com.agent.management.rag.kb.KnowledgeBaseRagRetriever; import com.agent.management.rag.kb.RagKnowledgeBaseConfigService; import com.agent.management.rag.model.*; import com.agent.management.repository.RagKnowledgeBaseRepository; import lombok.RequiredArgsConstructor; import org.springframework.web.bind.annotation.*; import java.util.*;
+@RestController @RequestMapping("/api/rag/kb") @RequiredArgsConstructor public class KnowledgeBaseRagDebugController {private final KnowledgeBaseRagRetriever retriever;private final RagKnowledgeBaseConfigService configs;private final RagKnowledgeBaseRepository knowledgeBaseRepository;@PostMapping("/retrieve") public Result<KnowledgeBaseRagResult> retrieve(@RequestBody KnowledgeBaseRagRequest q){return Result.success(retriever.retrieve(q));}@GetMapping("/{id}/bindings") public Result<List<RagKnowledgeSourceBinding>> bindings(@PathVariable Long id){return Result.success(configs.listBindings(id));}@PutMapping("/{id}/bindings/{bindingId}/config") public Result<RagKnowledgeSourceBinding> updateConfig(@PathVariable Long id,@PathVariable Long bindingId,@RequestBody Map<String,Object> config){return Result.success(configs.updateConfig(id,bindingId,config));}@PutMapping("/{id}/bindings/by-source") public Result<RagKnowledgeSourceBinding> upsertBySource(@PathVariable Long id,@RequestBody Map<String,Object> body){RagSourceType sourceType=RagSourceType.valueOf(String.valueOf(body.get("sourceType")).toUpperCase(Locale.ROOT));String sourceId=String.valueOf(body.get("sourceId"));@SuppressWarnings("unchecked") Map<String,Object> config=(Map<String,Object>)body.get("config");return Result.success(configs.upsertBySource(id,sourceType,sourceId,config));}@GetMapping public Result<List<RagKnowledgeBase>> list(){return Result.success(knowledgeBaseRepository.findAll());}}

+ 35 - 2
backend/src/main/java/com/agent/management/rag/kb/RagKnowledgeBaseConfigService.java

@@ -1,6 +1,7 @@
 package com.agent.management.rag.kb;
 
 import com.agent.management.model.entity.RagKnowledgeSourceBinding;
+import com.agent.management.rag.model.RagSourceType;
 import com.agent.management.repository.RagKnowledgeSourceBindingRepository;
 import com.fasterxml.jackson.databind.ObjectMapper;
 import lombok.RequiredArgsConstructor;
@@ -32,10 +33,42 @@ public class RagKnowledgeBaseConfigService {
         }
     }
 
+    /**
+     * 按 (knowledgeBaseId, sourceType, sourceId) 查找绑定:
+     * 存在则更新 configJson;不存在则新建绑定后再写入 configJson。
+     * 用于 RAG 治理页"应用审核后的授权配置":当用户尚未在知识库与数据源之间建立绑定时,
+     * 自动创建一条绑定记录,避免"没有匹配的数据源绑定"错误阻断治理流程。
+     */
+    public RagKnowledgeSourceBinding upsertBySource(Long knowledgeBaseId, RagSourceType sourceType,
+                                                     String sourceId, Map<String, Object> config) {
+        if (sourceType == null) {
+            throw new IllegalArgumentException("sourceType is required");
+        }
+        if (sourceId == null || sourceId.isBlank()) {
+            throw new IllegalArgumentException("sourceId is required");
+        }
+        RagKnowledgeSourceBinding binding = bindings
+                .findByKnowledgeBaseIdAndSourceTypeAndSourceId(knowledgeBaseId, sourceType, sourceId)
+                .orElseGet(() -> {
+                    RagKnowledgeSourceBinding created = new RagKnowledgeSourceBinding();
+                    created.setKnowledgeBaseId(knowledgeBaseId);
+                    created.setSourceType(sourceType);
+                    created.setSourceId(sourceId);
+                    return created;
+                });
+        validateAuthorization(binding, config);
+        try {
+            binding.setConfigJson(mapper.writeValueAsString(config == null ? Map.of() : config));
+            return bindings.save(binding);
+        } catch (Exception e) {
+            throw new IllegalArgumentException("invalid binding config: " + e.getMessage(), e);
+        }
+    }
+
     private static void validateAuthorization(RagKnowledgeSourceBinding binding, Map<String,Object> config) {
         if (config == null || !"AUTO_GENERATE".equals(String.valueOf(config.get("retrievalMode")))) return;
-        String key = binding.getSourceType() == com.agent.management.rag.model.RagSourceType.GRAPH ? "allowedLabels"
-                : binding.getSourceType() == com.agent.management.rag.model.RagSourceType.STRUCTURED_DATA ? "allowedTables" : null;
+        String key = binding.getSourceType() == RagSourceType.GRAPH ? "allowedLabels"
+                : binding.getSourceType() == RagSourceType.STRUCTURED_DATA ? "allowedTables" : null;
         if (key != null && (!(config.get(key) instanceof List<?> values) || values.isEmpty()))
             throw new IllegalArgumentException("automatic generation requires a non-empty explicit " + key);
     }

+ 2 - 2
backend/src/main/java/com/agent/management/repository/RagKnowledgeSourceBindingRepository.java

@@ -1,2 +1,2 @@
-package com.agent.management.repository; import com.agent.management.model.entity.RagKnowledgeSourceBinding; import org.springframework.data.jpa.repository.JpaRepository; import java.util.*;
-public interface RagKnowledgeSourceBindingRepository extends JpaRepository<RagKnowledgeSourceBinding,Long>{List<RagKnowledgeSourceBinding> findByKnowledgeBaseIdAndEnabledTrueOrderByPriorityAsc(Long knowledgeBaseId);List<RagKnowledgeSourceBinding> findByKnowledgeBaseIdOrderByPriorityAsc(Long knowledgeBaseId);}
+package com.agent.management.repository; import com.agent.management.model.entity.RagKnowledgeSourceBinding; import com.agent.management.rag.model.RagSourceType; import org.springframework.data.jpa.repository.JpaRepository; import java.util.*;
+public interface RagKnowledgeSourceBindingRepository extends JpaRepository<RagKnowledgeSourceBinding,Long>{List<RagKnowledgeSourceBinding> findByKnowledgeBaseIdAndEnabledTrueOrderByPriorityAsc(Long knowledgeBaseId);List<RagKnowledgeSourceBinding> findByKnowledgeBaseIdOrderByPriorityAsc(Long knowledgeBaseId);Optional<RagKnowledgeSourceBinding> findByKnowledgeBaseIdAndSourceTypeAndSourceId(Long knowledgeBaseId,RagSourceType sourceType,String sourceId);}

+ 7 - 1
backend/src/main/resources/application.yml.example

@@ -4,8 +4,14 @@ server:
 rag:
   ai-bridge:
     enabled: ${RAG_AI_BRIDGE_ENABLED:false}
-    base-url: ${RAG_AI_BRIDGE_BASE_URL:http://127.0.0.1:18733}
     timeout: ${RAG_AI_BRIDGE_TIMEOUT:300s}
+    # Python 子进程管理(由 Java 启动并注入 LLM 环境变量,LLM 配置统一从 spring.ai.openai.* 读取)
+    host: ${RAG_AI_BRIDGE_HOST:127.0.0.1}
+    port: ${RAG_AI_BRIDGE_PORT:18733}
+    python-path: ${RAG_AI_BRIDGE_PYTHON_PATH:python}
+    script-path: ${RAG_AI_BRIDGE_SCRIPT_PATH:rag-ai-bridge/server.py}
+    startup-timeout: ${RAG_AI_BRIDGE_STARTUP_TIMEOUT:60}
+    health-check-interval: ${RAG_AI_BRIDGE_HEALTH_CHECK_INTERVAL:60}
 
 skill:
   # Skill 文件存放目录,请修改为你的实际路径

+ 8 - 0
frontend/src/api/rag.js

@@ -56,10 +56,18 @@ export function getKnowledgeBaseBindings(knowledgeBaseId) {
   return unwrap(request.get(`/rag/kb/${knowledgeBaseId}/bindings`))
 }
 
+export function listKnowledgeBases() {
+  return unwrap(request.get('/rag/kb'))
+}
+
 export function updateKnowledgeBaseBindingConfig(knowledgeBaseId, bindingId, config) {
   return unwrap(request.put(`/rag/kb/${knowledgeBaseId}/bindings/${bindingId}/config`, config))
 }
 
+export function upsertKnowledgeBaseBindingBySource(knowledgeBaseId, sourceType, sourceId, config) {
+  return unwrap(request.put(`/rag/kb/${knowledgeBaseId}/bindings/by-source`, { sourceType, sourceId, config }))
+}
+
 export function generateRagAnswer(question, evidences) {
   return unwrap(request.post('/rag/answer', { question, evidences }, { timeout: RAG_TIMEOUT }))
 }

+ 69 - 0
frontend/src/components/workflow/nodes/KnowledgeRetrievalNode.vue

@@ -0,0 +1,69 @@
+<script setup>
+import { computed } from 'vue'
+import BaseNode from './BaseNode.vue'
+import { LibraryOutline } from '@vicons/ionicons5'
+
+const props = defineProps({
+  data: { type: Object, default: () => ({ label: '知识库检索', source: 'document' }) }
+})
+
+const SOURCE_LABELS = {
+  document: '文档数据库',
+  structured: '结构化数据库',
+  graph: '知识图谱库',
+  hybrid: '混合检索'
+}
+
+const sourceLabel = computed(() => SOURCE_LABELS[props.data?.source] || '文档数据库')
+
+const querySummary = computed(() => {
+  const q = props.data?.query || ''
+  if (!q) return '未填写检索语句'
+  return q.length > 30 ? q.slice(0, 30) + '…' : q
+})
+
+const inputSummary = computed(() => (props.data?.inputs || []).map(v => v.name).filter(Boolean))
+const outputSummary = computed(() => (props.data?.outputs || []).map(v => v.name).filter(Boolean))
+</script>
+
+<template>
+  <BaseNode :data="data" color="#10B981">
+    <template #icon><LibraryOutline /></template>
+    <template #body>
+      <div class="kr-source">{{ sourceLabel }}</div>
+      <div class="kr-query">{{ querySummary }}</div>
+      <div v-if="inputSummary.length" class="io-summary">
+        <span class="io-tag io-in" v-for="name in inputSummary" :key="name">{{ name }}</span>
+      </div>
+      <div v-if="outputSummary.length" class="io-summary">
+        <span class="io-tag io-out" v-for="name in outputSummary" :key="name">{{ name }}</span>
+      </div>
+    </template>
+  </BaseNode>
+</template>
+
+<script>
+export default { name: 'KnowledgeRetrievalNode' }
+</script>
+
+<style scoped>
+.kr-source { font-size: 11px; color: var(--text-tertiary, #888); margin-bottom: 2px; }
+.kr-query {
+  font-size: 11px;
+  color: var(--text-secondary, #aaa);
+  margin-bottom: 4px;
+  max-width: 220px;
+  overflow: hidden;
+  text-overflow: ellipsis;
+  white-space: nowrap;
+}
+.io-summary { display: flex; flex-wrap: wrap; gap: 3px; margin-top: 3px; }
+.io-tag {
+  font-size: 10px;
+  padding: 1px 5px;
+  border-radius: 3px;
+  font-family: 'Cascadia Code', 'Fira Code', monospace;
+}
+.io-in { background: rgba(99,102,241,0.15); color: #818cf8; }
+.io-out { background: rgba(34,197,94,0.15); color: #4ade80; }
+</style>

+ 15 - 0
frontend/src/utils/ioInference.js

@@ -68,6 +68,14 @@ export function getNodeOutputs(node) {
       return d.outputs || []
     case 'skill':
       return d.outputs || []
+    case 'knowledgeRetrieval':
+      if (d.outputs && d.outputs.length) return d.outputs
+      return [
+        { name: 'evidences', label: '检索证据列表', type: 'array', description: '命中的知识片段集合' },
+        { name: 'evidenceCount', label: '证据数量', type: 'number', description: '命中证据条数' },
+        { name: 'sourceType', label: '来源类型', type: 'string', description: 'DOCUMENT / STRUCTURED_DATA / GRAPH / HYBRID' },
+        { name: 'diagnostics', label: '诊断信息', type: 'object', description: '检索过程的诊断信息' }
+      ]
     case 'output':
       return []
     case 'condition':
@@ -100,6 +108,13 @@ export function getNodeInputs(node) {
       return d.inputs || []
     case 'skill':
       return d.inputs || []
+    case 'knowledgeRetrieval': {
+      // 优先使用手动定义的 inputs
+      if (d.inputs && d.inputs.length) return d.inputs
+      // 回退:从检索语句模板提取
+      const krVars = new Set(extractTemplateVariables(d.query))
+      return [...krVars].map(name => ({ name, label: name, type: 'string', description: '' }))
+    }
     case 'output':
       return d.fields || []
     case 'condition': {

+ 3 - 3
frontend/src/views/knowledge/RagGovernance.vue

@@ -3,7 +3,7 @@ import { computed, onMounted, reactive, ref } from 'vue'
 import { useMessage } from 'naive-ui'
 import {
   addRagQueryExampleCandidates, deleteRagQueryExample, getKnowledgeBaseBindings, listRagQueryExamples, refreshRagCapability,
-  suggestRagGovernance, updateKnowledgeBaseBindingConfig, verifyRagQueryExample
+  suggestRagGovernance, upsertKnowledgeBaseBindingBySource, verifyRagQueryExample
 } from '../../api/rag'
 
 const message = useMessage()
@@ -147,7 +147,6 @@ function buildConfig() {
 }
 
 async function applyConfig() {
-  if (!binding.value) return message.error('知识库 1 中没有匹配的数据源绑定')
   if (isGraph.value && !selectedLabels.value.length) return message.error('至少授权一个 Label,空白名单会造成权限语义不明确')
   if (!isGraph.value && !selectedTables.value.length) return message.error('至少授权一个 Table,空白名单会造成权限语义不明确')
   if (isGraph.value) {
@@ -157,7 +156,8 @@ async function applyConfig() {
   }
   loading.apply = true
   try {
-    await updateKnowledgeBaseBindingConfig(1, binding.value.id, buildConfig())
+    // 按 (sourceType, sourceId) 查找或创建绑定:解决用户进入治理页时 kbId=1 下尚未建立绑定记录的问题
+    await upsertKnowledgeBaseBindingBySource(1, form.sourceType, form.sourceId, buildConfig())
     bindings.value = await getKnowledgeBaseBindings(1)
     message.success('授权配置已应用')
   } catch (error) { message.error(error.message || '配置应用失败') }

+ 241 - 2
frontend/src/views/workflow/WorkflowEditor.vue

@@ -4,23 +4,28 @@ import { useRoute, useRouter } from 'vue-router'
 import { VueFlow, useVueFlow } from '@vue-flow/core'
 import { Background } from '@vue-flow/background'
 import {
-  NInput, NSelect, NButton, NIcon, useMessage
+  NInput, NSelect, NButton, NIcon, NCheckbox, useMessage
 } from 'naive-ui'
 import {
   ArrowBackOutline, SaveOutline, ChatbubbleEllipsesOutline,
   SparklesOutline, RocketOutline, FlaskOutline, GitBranchOutline, ExitOutline,
   PlayOutline, CloseOutline, ColorWandOutline, LinkOutline, TrashOutline, OpenOutline,
-  TimeOutline
+  TimeOutline, LibraryOutline
 } from '@vicons/ionicons5'
 import { useWorkflowStore } from '../../stores/workflow'
 import { useSkillStore } from '../../stores/skill'
 import { runWorkflow, autoAssociate } from '../../api/workflow'
 import { getAgentTemplate } from '../../api/agentTemplate'
 import { getModelList } from '../../api/model'
+import { getDataSources } from '../../api/datasource'
+import { getGraphSources } from '../../api/graphsource'
+import { getKbDocuments } from '../../api/knowledge'
+import { listKnowledgeBases } from '../../api/rag'
 import InputNode from '../../components/workflow/nodes/InputNode.vue'
 import LLMNode from '../../components/workflow/nodes/LLMNode.vue'
 import AgentNode from '../../components/workflow/nodes/AgentNode.vue'
 import SkillNode from '../../components/workflow/nodes/SkillNode.vue'
+import KnowledgeRetrievalNode from '../../components/workflow/nodes/KnowledgeRetrievalNode.vue'
 import OutputNode from '../../components/workflow/nodes/OutputNode.vue'
 import ConditionNode from '../../components/workflow/nodes/ConditionNode.vue'
 import SmartActionNode from '../../components/workflow/nodes/SmartActionNode.vue'
@@ -106,6 +111,7 @@ const {
     agent:       { color: '#F59E0B', typeLabel: '智能体' },
     smartAction: { color: '#EC4899', typeLabel: '智能操作' },
     condition:   { color: '#F97316', typeLabel: '条件' },
+    knowledgeRetrieval: { color: '#10B981', typeLabel: '知识库检索' },
     output:      { color: '#EF4444', typeLabel: '输出' }
   },
   workflowIdRef: workflowId
@@ -192,6 +198,7 @@ const nodeTypes = [
   { type: 'skill', label: '技能', icon: markRaw(FlaskOutline), color: '#06B6D4' },
   { type: 'agent', label: '智能体', icon: markRaw(RocketOutline), color: '#F59E0B' },
   { type: 'smartAction', label: '智能操作', icon: markRaw(ColorWandOutline), color: '#EC4899' },
+  { type: 'knowledgeRetrieval', label: '知识库检索', icon: markRaw(LibraryOutline), color: '#10B981' },
   { type: 'condition', label: '条件分支', icon: markRaw(GitBranchOutline), color: '#F97316' },
   { type: 'output', label: '输出', icon: markRaw(ExitOutline), color: '#EF4444' }
 ]
@@ -202,6 +209,7 @@ const nodeTypeInfo = {
   skill:       { color: '#06B6D4', typeLabel: '技能',     icon: markRaw(FlaskOutline) },
   agent:       { color: '#F59E0B', typeLabel: '智能体',   icon: markRaw(RocketOutline) },
   smartAction: { color: '#EC4899', typeLabel: '智能操作', icon: markRaw(ColorWandOutline) },
+  knowledgeRetrieval: { color: '#10B981', typeLabel: '知识库检索', icon: markRaw(LibraryOutline) },
   condition:   { color: '#F97316', typeLabel: '条件',     icon: markRaw(GitBranchOutline) },
   output:      { color: '#EF4444', typeLabel: '输出',     icon: markRaw(ExitOutline) }
 }
@@ -213,6 +221,19 @@ function getDefaultData(type) {
     case 'agent': return { label: '智能体', agentId: '', agentName: '', modelId: null, failStrategy: 'abort' }
     case 'skill': return { label: '技能', skillId: '', skillName: '', modelId: null, failStrategy: 'abort' }
     case 'smartAction': return { label: '智能操作', actionPrompt: '', modelId: null, inputs: [], outputs: [{ name: 'result', label: '操作结果', type: 'string', description: '' }], failStrategy: 'abort' }
+    case 'knowledgeRetrieval': return {
+      label: '知识库检索',
+      source: 'document',
+      query: '',
+      allDocuments: true,
+      documentIds: [],
+      allDatasources: false,
+      datasourceIds: [],
+      allGraphSources: false,
+      graphSourceIds: [],
+      knowledgeBaseId: null,
+      topK: 5
+    }
     case 'condition': return { label: '条件分支', conditions: [{ type: 'IF', expression: '' }] }
     case 'output': return { label: '输出', fields: [{ name: 'result', label: '输出结果', type: 'string', description: '' }] }
     default: return { label: '节点' }
@@ -369,6 +390,83 @@ const failStrategyOptions = [
   { label: '跳过继续', value: 'skip' }
 ]
 
+// ========== 知识库检索节点:下拉/复选框数据源 ==========
+const krSourceOptions = [
+  { label: '文档数据库', value: 'document' },
+  { label: '结构化数据库', value: 'structured' },
+  { label: '知识图谱库', value: 'graph' },
+  { label: '混合检索', value: 'hybrid' }
+]
+const krDocuments = ref([])
+const krDocumentsLoading = ref(false)
+const krDatasources = ref([])
+const krDatasourcesLoading = ref(false)
+const krGraphSources = ref([])
+const krGraphSourcesLoading = ref(false)
+const krKnowledgeBases = ref([])
+const krKnowledgeBasesLoading = ref(false)
+
+const krDocumentOptions = computed(() => krDocuments.value.map(d => ({ label: d.name || `#${d.id}`, value: String(d.id) })))
+const krDatasourceOptions = computed(() => krDatasources.value.map(d => ({ label: d.name || `#${d.id}`, value: String(d.id) })))
+const krGraphSourceOptions = computed(() => krGraphSources.value.map(d => ({ label: d.name || `#${d.id}`, value: String(d.id) })))
+const krKnowledgeBaseOptions = computed(() => krKnowledgeBases.value.map(kb => ({ label: kb.name || `#${kb.id}`, value: kb.id })))
+
+async function loadKrDocuments() {
+  if (krDocuments.value.length || krDocumentsLoading.value) return
+  krDocumentsLoading.value = true
+  try {
+    const res = await getKbDocuments({ page: 1, size: 1000 })
+    krDocuments.value = res.data?.items || []
+  } catch (e) {
+    message.error('加载文档列表失败')
+  } finally {
+    krDocumentsLoading.value = false
+  }
+}
+async function loadKrDatasources() {
+  if (krDatasources.value.length || krDatasourcesLoading.value) return
+  krDatasourcesLoading.value = true
+  try {
+    const res = await getDataSources()
+    krDatasources.value = res.data || []
+  } catch (e) {
+    message.error('加载结构化数据源失败')
+  } finally {
+    krDatasourcesLoading.value = false
+  }
+}
+async function loadKrGraphSources() {
+  if (krGraphSources.value.length || krGraphSourcesLoading.value) return
+  krGraphSourcesLoading.value = true
+  try {
+    const res = await getGraphSources()
+    krGraphSources.value = res.data || []
+  } catch (e) {
+    message.error('加载图谱数据源失败')
+  } finally {
+    krGraphSourcesLoading.value = false
+  }
+}
+async function loadKrKnowledgeBases() {
+  if (krKnowledgeBases.value.length || krKnowledgeBasesLoading.value) return
+  krKnowledgeBasesLoading.value = true
+  try {
+    krKnowledgeBases.value = await listKnowledgeBases() || []
+  } catch (e) {
+    message.error('加载知识库列表失败')
+  } finally {
+    krKnowledgeBasesLoading.value = false
+  }
+}
+
+function toggleKrListField(field, value) {
+  if (!selectedData.value) return
+  const list = new Set(selectedData.value[field] || [])
+  if (list.has(value)) list.delete(value)
+  else list.add(value)
+  onFieldChange(field, [...list])
+}
+
 // ========== 节点字段操作(来自 composable) ==========
 const {
   onFieldChange, onModelChange,
@@ -753,6 +851,11 @@ async function saveMessageWrap(v) {
 onMounted(async () => {
   skillStore.fetchSkills()
   fetchDbModels()
+  // 预加载知识库检索节点所需的资源列表(失败不阻塞主流程)
+  loadKrDocuments()
+  loadKrDatasources()
+  loadKrGraphSources()
+  loadKrKnowledgeBases()
   if (workflowId.value) {
     // 编辑已有工作流
     try {
@@ -901,6 +1004,7 @@ function applyGraphData(graphData) {
           <template #node-llm="nodeProps"><LLMNode :data="nodeProps.data" /></template>
           <template #node-agent="nodeProps"><AgentNode :data="nodeProps.data" /></template>
           <template #node-skill="nodeProps"><SkillNode :data="nodeProps.data" /></template>
+          <template #node-knowledgeRetrieval="nodeProps"><KnowledgeRetrievalNode :data="nodeProps.data" /></template>
           <template #node-output="nodeProps"><OutputNode :data="nodeProps.data" /></template>
           <template #node-condition="nodeProps"><ConditionNode :data="nodeProps.data" /></template>
           <template #node-smartAction="nodeProps"><SmartActionNode :data="nodeProps.data" /></template>
@@ -1249,6 +1353,126 @@ function applyGraphData(graphData) {
                 </div>
               </template>
 
+              <!-- 知识库检索节点 -->
+              <template v-if="selectedNode?.type === 'knowledgeRetrieval'">
+                <div class="prop-section">
+                  <label class="prop-label">检索来源</label>
+                  <n-select
+                    :value="selectedData?.source || 'document'"
+                    :options="krSourceOptions"
+                    size="small"
+                    @update:value="v => onFieldChange('source', v)"
+                  />
+                </div>
+                <div class="prop-section">
+                  <label class="prop-label">检索语句</label>
+                  <n-input
+                    :value="selectedData?.query"
+                    type="textarea" :rows="3"
+                    placeholder="输入检索语句,可用 {{变量名}} 引用上游变量"
+                    size="small"
+                    @update:value="v => onFieldChange('query', v)"
+                  />
+                </div>
+
+                <!-- 文档数据库 -->
+                <template v-if="(selectedData?.source || 'document') === 'document'">
+                  <div class="prop-section">
+                    <label class="kr-checkbox-row">
+                      <n-checkbox
+                        :checked="!!selectedData?.allDocuments"
+                        @update:checked="v => onFieldChange('allDocuments', v)"
+                      />
+                      <span>在所有文档中检索</span>
+                    </label>
+                    <div v-if="!selectedData?.allDocuments" class="kr-list">
+                      <div v-if="krDocumentsLoading" class="empty-hint">加载中...</div>
+                      <div v-else-if="!krDocumentOptions.length" class="empty-hint">暂无文档</div>
+                      <label v-for="opt in krDocumentOptions" :key="opt.value" class="kr-checkbox-row">
+                        <n-checkbox
+                          :checked="(selectedData?.documentIds || []).includes(opt.value)"
+                          @update:checked="() => toggleKrListField('documentIds', opt.value)"
+                        />
+                        <span>{{ opt.label }}</span>
+                      </label>
+                    </div>
+                  </div>
+                  <div class="prop-section">
+                    <label class="prop-label">TopK(命中条数)</label>
+                    <n-input
+                      :value="String(selectedData?.topK ?? 5)"
+                      type="text"
+                      size="small"
+                      @update:value="v => onFieldChange('topK', Math.max(1, parseInt(v) || 5))"
+                    />
+                  </div>
+                </template>
+
+                <!-- 结构化数据库 -->
+                <template v-if="selectedData?.source === 'structured'">
+                  <div class="prop-section">
+                    <label class="kr-checkbox-row">
+                      <n-checkbox
+                        :checked="!!selectedData?.allDatasources"
+                        @update:checked="v => onFieldChange('allDatasources', v)"
+                      />
+                      <span>在所有数据源中检索</span>
+                    </label>
+                    <div v-if="!selectedData?.allDatasources" class="kr-list">
+                      <div v-if="krDatasourcesLoading" class="empty-hint">加载中...</div>
+                      <div v-else-if="!krDatasourceOptions.length" class="empty-hint">暂无数据源</div>
+                      <label v-for="opt in krDatasourceOptions" :key="opt.value" class="kr-checkbox-row">
+                        <n-checkbox
+                          :checked="(selectedData?.datasourceIds || []).includes(opt.value)"
+                          @update:checked="() => toggleKrListField('datasourceIds', opt.value)"
+                        />
+                        <span>{{ opt.label }}</span>
+                      </label>
+                    </div>
+                  </div>
+                </template>
+
+                <!-- 知识图谱库 -->
+                <template v-if="selectedData?.source === 'graph'">
+                  <div class="prop-section">
+                    <label class="kr-checkbox-row">
+                      <n-checkbox
+                        :checked="!!selectedData?.allGraphSources"
+                        @update:checked="v => onFieldChange('allGraphSources', v)"
+                      />
+                      <span>在所有数据源中检索</span>
+                    </label>
+                    <div v-if="!selectedData?.allGraphSources" class="kr-list">
+                      <div v-if="krGraphSourcesLoading" class="empty-hint">加载中...</div>
+                      <div v-else-if="!krGraphSourceOptions.length" class="empty-hint">暂无图谱数据源</div>
+                      <label v-for="opt in krGraphSourceOptions" :key="opt.value" class="kr-checkbox-row">
+                        <n-checkbox
+                          :checked="(selectedData?.graphSourceIds || []).includes(opt.value)"
+                          @update:checked="() => toggleKrListField('graphSourceIds', opt.value)"
+                        />
+                        <span>{{ opt.label }}</span>
+                      </label>
+                    </div>
+                  </div>
+                </template>
+
+                <!-- 混合检索 -->
+                <template v-if="selectedData?.source === 'hybrid'">
+                  <div class="prop-section">
+                    <label class="prop-label">知识库</label>
+                    <div v-if="krKnowledgeBasesLoading" class="empty-hint">加载中...</div>
+                    <n-select
+                      :value="selectedData?.knowledgeBaseId"
+                      :options="krKnowledgeBaseOptions"
+                      placeholder="选择知识库"
+                      size="small"
+                      @update:value="v => onFieldChange('knowledgeBaseId', v)"
+                    />
+                    <div v-if="!krKnowledgeBaseOptions.length && !krKnowledgeBasesLoading" class="empty-hint">暂无知识库</div>
+                  </div>
+                </template>
+              </template>
+
               <!-- 条件分支节点 -->
               <template v-if="selectedNode?.type === 'condition'">
                 <div class="prop-section">
@@ -2127,6 +2351,21 @@ function applyGraphData(graphData) {
   font-style: italic;
 }
 
+.kr-checkbox-row {
+  display: flex;
+  align-items: center;
+  gap: 6px;
+  font-size: 12px;
+  cursor: pointer;
+  margin-bottom: 4px;
+}
+.kr-list {
+  margin-top: 6px;
+  max-height: 240px;
+  overflow-y: auto;
+  padding-right: 4px;
+}
+
 .passthrough-hint {
   font-size: 11px;
   color: var(--text-tertiary, #555);

+ 84 - 0
prompt.md

@@ -985,3 +985,87 @@ P0方案B,先看看java的技能为什么不被hermes识别
 
 ---
 
+@agent-management-rag目录,是我针对该项目的另一个worktree,主要做了RAG的检索和治理的相关工作。请分析提交记录,然后合并到主分支。
+
+---
+
+丢弃暂存区更改和当前文件更改,还原到最近一次提交,生成命令;删除所有Todo
+
+---
+
+读取@agent-management-rag/02_907f73440ae2_git_diff.patch文件,应用其中的更改。注意:
+1. 该更改是其他同事的更改,而我已经有了大量其他更改。所以,行号有可能已变化。所以,不要使用git命令进行合并,而是你来读取文件内容,然后智能把更改写入当前分支。尽量不要使用脚本来进行patch应用。
+2. 不要读取其他patch文件。
+3. 对于新增文件,直接复制到对应路径即可;对于修改文件,由你来进行智能修改。
+
+---
+
+@temp/co-defense-rag-call-examples.md中,是针对某个场景进行知识库检索的调用示例。
+参考@temp/co-defense-rag-call-examples.md中调用知识库检索的示例,在工作流中增加节点“知识库检索”。右侧属性配置界面,可下拉选择4种检索来源:文档数据库、结构化数据库、知识图谱库、混合检索。
+1. 选择文档数据库时,生成以下表单:
+  - 检索语句(使用自然语言检索,多行文本框,例如“协防关系定义规则是什么?”)
+  - 在所有文档中检索(复选框,选择后,下方“文档范围”条目消失)
+  - 文档范围(列出当前文档数据库中的文档,用复选框选择,后续只在选择的文档中搜索)
+  - TopK(大于0的数字)
+2. 选择结构化数据库时,生成以下表单:
+  - 检索语句(使用自然语言检索,多行文本框,例如“查询XX航母的具体属性”)
+  - 在所有数据源中检索(复选框,选择后,下方“生效数据源”条目消失)
+  - 生效数据源(列出当前所有的结构化数据源,用复选框选择,后续只在选择的数据源中搜索)
+3. 选择知识图谱库时,生成以下表单:
+  - 检索语句(使用自然语言检索,多行文本框,例如“查询与XX航母有协防关系的装备信息”)
+  - 在所有数据源中检索(复选框,选择后,下方“生效数据源”条目消失)
+  - 生效数据源(列出当前所有的知识图谱数据源,用复选框选择,后续只在选择的数据源中搜索)
+
+---
+
+执行文档检索;报错:
+
+**注意:不要搜索,请到相关的py文件中定位问题;问题不在于部署的是milvus还是milvus lite。**
+
+---
+
+删除所有TODO。我将milvus-lite和pymilvus版本统一了,都是3.0,然后重启服务,报错:
+
+**注意:不要搜索,请到相关的py文件中定位问题;问题不在于部署的是milvus还是milvus lite。**
+
+---
+
+文档数据可以搜索出来了,但RAG页面,混合检索报错:
+ERROR c.a.m.c.exception.GlobalExceptionHandler - 系统异常
+java.lang.IllegalStateException: rag-ai-bridge is disabled
+
+---
+
+我希望只在application.yml里配置一次大模型,rag-ai-bridge启动时,Java通过环境变量等注入python。
+
+---
+
+图谱检索,报错:
+{
+  "graph:1": {
+    "error": "Cypher generation failed: 503 Service Unavailable: \"{\"detail\":\"query repair unavailable: Request timed out.\"}\""
+  }
+}
+
+---
+
+模型能力没问题。现在仍有报错:
+{
+  "graph:1": {
+    "error": "Cypher generation failed: Cypher uses nonexistent or unauthorized label: POTENTIAL_SUPPORT"
+  }
+}
+
+---
+
+将工作流执行超时时间,以及每个节点超时时间设为无限。
+
+---
+
+当前,工作流有两个问题:
+1. 智能操作、技能、大模型等节点在执行时,看不到前面节点的输出。我需要你将前置节点输出的工作空间各变量作为系统提示词注入。
+2. 如果某节点后置多个节点,不会并行执行,而是串行挨个执行。
+
+---
+
+在当前文件的逻辑中,点击“应用审核后的授权配置”按钮,提示“知识库 1 中没有匹配的数据源绑定”。看起来没有和我真实的数据源关联起来。请解决。