feat: 增加文件上传相关逻辑

This commit is contained in:
qiufeng
2026-07-23 10:45:13 +08:00
parent 10ce992679
commit ead8b66bd3
32 changed files with 3093 additions and 74 deletions
+176 -1
View File
@@ -34,6 +34,181 @@
## 🤖AI能力
<style>
:root {
--bg: #ffffff;
--surface: #ffffff;
--surface-muted: #f6f6f8;
--border: #e4e4e7;
--border-strong: #d4d4d8;
--text: #18181b;
--text-muted: #71717a;
--brand: #7c3aed;
--brand-bg: #f3eafe;
--brand-border: #c4b5fd;
--external: #b45309;
--external-bg: #fef3c7;
--external-border: #fcd34d;
--connector: #a1a1aa;
--font-sans: -apple-system, BlinkMacSystemFont, 'Segoe UI', Roboto, sans-serif;
--font-mono: 'SF Mono', Monaco, Consolas, monospace;
--weight-medium: 500;
--weight-strong: 600;
}
:root[data-widget-theme="dark"] {
--bg: #0f0f12;
--surface: #18181b;
--surface-muted: #27272a;
--border: #3f3f46;
--border-strong: #52525b;
--text: #f4f4f5;
--text-muted: #a1a1aa;
--brand: #a78bfa;
--brand-bg: #2e1065;
--brand-border: #6d28d9;
--external: #fbbf24;
--external-bg: #422006;
--external-border: #92400e;
--connector: #71717a;
}
.arch-card {
background: var(--surface);
border: 1px solid var(--border);
border-radius: 12px;
padding: 16px;
}
.arch-svg { width: 100%; height: auto; display: block; }
.arch-svg text { font-family: var(--font-sans); }
.arch-svg .t-label { font-size: 14px; font-weight: var(--weight-medium); fill: var(--text); }
.arch-svg .t-sub { font-size: 12px; fill: var(--text-muted); }
.arch-svg .t-title { font-size: 13px; font-weight: var(--weight-strong); fill: var(--text-muted); letter-spacing: 0.3px; }
.arch-svg .t-mini { font-size: 11px; fill: var(--text-muted); }
.arch-svg .t-brand { fill: var(--brand); font-weight: var(--weight-medium); }
.arch-svg .node-neutral { fill: var(--surface); stroke: var(--border-strong); stroke-width: 1; }
.arch-svg .node-brand { fill: var(--brand-bg); stroke: var(--brand); stroke-width: 1.5; }
.arch-svg .node-external { fill: var(--external-bg); stroke: var(--external-border); stroke-width: 1; stroke-dasharray: 4 3; }
.arch-svg .boundary { fill: none; stroke: var(--border); stroke-width: 1; stroke-dasharray: 5 4; }
.arch-svg .conn { fill: none; stroke: var(--connector); stroke-width: 1.4; stroke-linecap: round; stroke-linejoin: round; }
.arch-svg .conn-brand { fill: none; stroke: var(--brand); stroke-width: 1.8; stroke-linecap: round; stroke-linejoin: round; }
.arch-legend { display: flex; gap: 20px; flex-wrap: wrap; margin-top: 12px; font-size: 12px; color: var(--text-muted); font-family: var(--font-sans); }
.arch-legend span { display: inline-flex; align-items: center; gap: 6px; }
.arch-legend .sw { width: 18px; height: 10px; border-radius: 3px; display: inline-block; }
.arch-note { margin-top: 10px; font-size: 12px; color: var(--text-muted); font-family: var(--font-sans); line-height: 1.5; }
</style>
<div class="arch-card" data-dynamic-ui-widget data-template="node-flow" data-mounted="true">
<svg class="arch-svg" viewBox="0 0 720 810" role="img" aria-label="mfish AI 能力架构图">
<defs>
<marker id="arr" viewBox="0 0 8 8" refX="7" refY="4" markerWidth="8" markerHeight="8" markerUnits="userSpaceOnUse" orient="auto">
<path d="M1 1 L7 4 L1 7 Z" fill="#a1a1aa"/>
</marker>
<marker id="arr-brand" viewBox="0 0 8 8" refX="7" refY="4" markerWidth="8" markerHeight="8" markerUnits="userSpaceOnUse" orient="auto">
<path d="M1 1 L7 4 L1 7 Z" fill="#7c3aed"/>
</marker>
</defs>
<!-- 前端 -->
<rect x="300" y="24" width="120" height="38" rx="19" class="node-neutral"/>
<text x="360" y="48" text-anchor="middle" class="t-label">前端 Frontend</text>
<!-- 入口层 -->
<rect x="40" y="86" width="640" height="118" rx="12" class="boundary"/>
<text x="56" y="104" class="t-title">入口层 · Controllers</text>
<rect x="65" y="116" width="180" height="36" rx="8" class="node-neutral"/>
<text x="155" y="139" text-anchor="middle" class="t-label">AiController</text>
<rect x="270" y="116" width="180" height="36" rx="8" class="node-neutral"/>
<text x="360" y="139" text-anchor="middle" class="t-label">RouterController</text>
<rect x="475" y="116" width="180" height="36" rx="8" class="node-neutral"/>
<text x="565" y="139" text-anchor="middle" class="t-label">AgentController</text>
<rect x="270" y="162" width="180" height="30" rx="8" class="node-neutral"/>
<text x="360" y="181" text-anchor="middle" class="t-sub">AiRouteService · LLM路由</text>
<!-- 助手层 -->
<rect x="40" y="226" width="640" height="160" rx="12" class="boundary"/>
<text x="56" y="244" class="t-title">助手层 · Assistants</text>
<!-- BaseAssistant 分组 -->
<rect x="60" y="258" width="320" height="120" rx="8" class="boundary"/>
<rect x="80" y="272" width="280" height="32" rx="8" class="node-brand"/>
<text x="220" y="293" text-anchor="middle" class="t-label t-brand">BaseAssistant</text>
<text x="220" y="324" text-anchor="middle" class="t-mini">Mfish · Sys · Storage · Oauth</text>
<text x="220" y="342" text-anchor="middle" class="t-mini">Scheduler · NoCode · Workflow</text>
<text x="220" y="360" text-anchor="middle" class="t-mini">Gateway · Agent 9 子类)</text>
<!-- AgentRuntime 分组 -->
<rect x="400" y="258" width="280" height="120" rx="8" class="boundary"/>
<rect x="420" y="272" width="240" height="32" rx="8" class="node-brand"/>
<text x="540" y="293" text-anchor="middle" class="t-label t-brand">AgentRuntime</text>
<rect x="420" y="316" width="110" height="28" rx="8" class="node-neutral"/>
<text x="475" y="334" text-anchor="middle" class="t-sub">Planner</text>
<rect x="550" y="316" width="110" height="28" rx="8" class="node-neutral"/>
<text x="605" y="334" text-anchor="middle" class="t-sub">Executor</text>
<rect x="420" y="352" width="240" height="22" rx="8" class="node-neutral"/>
<text x="540" y="367" text-anchor="middle" class="t-mini">EventBus · SSE 事件流</text>
<!-- 能力层 -->
<rect x="40" y="406" width="640" height="110" rx="12" class="boundary"/>
<text x="56" y="424" class="t-title">能力层 · mf-ai Services</text>
<rect x="60" y="438" width="140" height="44" rx="8" class="node-neutral"/>
<text x="130" y="458" text-anchor="middle" class="t-label">FileParseService</text>
<text x="130" y="474" text-anchor="middle" class="t-mini">文件解析 · 新增</text>
<rect x="220" y="438" width="140" height="44" rx="8" class="node-neutral"/>
<text x="290" y="458" text-anchor="middle" class="t-label">LlmModelRouter</text>
<text x="290" y="474" text-anchor="middle" class="t-mini">按租户路由模型</text>
<rect x="380" y="438" width="140" height="44" rx="8" class="node-brand"/>
<text x="450" y="458" text-anchor="middle" class="t-label t-brand">ApiToolEngine</text>
<text x="450" y="474" text-anchor="middle" class="t-mini">工具聚合引擎</text>
<rect x="540" y="438" width="120" height="44" rx="8" class="node-neutral"/>
<text x="600" y="458" text-anchor="middle" class="t-sub">FaultTolerant</text>
<text x="600" y="474" text-anchor="middle" class="t-mini">ToolCallingManager</text>
<!-- 工具来源层 -->
<rect x="40" y="536" width="640" height="86" rx="12" class="boundary"/>
<text x="56" y="554" class="t-title">工具来源 · mf-common-ai · ToolProvider</text>
<rect x="100" y="568" width="240" height="38" rx="8" class="node-neutral"/>
<text x="220" y="585" text-anchor="middle" class="t-sub">FeignToolProvider</text>
<text x="220" y="599" text-anchor="middle" class="t-mini">扫描 @FeignClient 接口</text>
<rect x="380" y="568" width="240" height="38" rx="8" class="node-neutral"/>
<text x="500" y="585" text-anchor="middle" class="t-sub">HttpToolProvider</text>
<text x="500" y="599" text-anchor="middle" class="t-mini">基于 OpenAPI 文档</text>
<!-- 外部服务 -->
<rect x="70" y="642" width="180" height="44" rx="8" class="node-external"/>
<text x="160" y="662" text-anchor="middle" class="t-label">mf-storage</text>
<text x="160" y="678" text-anchor="middle" class="t-mini">文件存储服务</text>
<rect x="270" y="642" width="180" height="44" rx="8" class="node-external"/>
<text x="360" y="662" text-anchor="middle" class="t-label">微服务集群</text>
<text x="360" y="678" text-anchor="middle" class="t-mini">sys/oauth/scheduler/...</text>
<rect x="470" y="642" width="180" height="44" rx="8" class="node-external"/>
<text x="560" y="662" text-anchor="middle" class="t-label">LLM 提供商</text>
<text x="560" y="678" text-anchor="middle" class="t-mini">ChatModel 流式</text>
<!-- 连接线 -->
<path class="conn" d="M360 62 L360 86" marker-end="url(#arr)"/>
<path class="conn" d="M310 192 Q260 210 250 258" marker-end="url(#arr)"/>
<path class="conn" d="M410 192 Q470 210 520 258" marker-end="url(#arr)"/>
<path class="conn" d="M180 378 Q150 400 140 438" marker-end="url(#arr)"/>
<path class="conn-brand" d="M280 378 Q380 410 420 438" marker-end="url(#arr-brand)"/>
<path class="conn-brand" d="M540 378 Q500 410 470 438" marker-end="url(#arr-brand)"/>
<path class="conn-brand" d="M450 482 L450 568" marker-end="url(#arr-brand)"/>
<path class="conn" d="M130 482 Q140 560 150 642" marker-end="url(#arr)"/>
<path class="conn" d="M220 606 Q300 620 340 642" marker-end="url(#arr)"/>
<path class="conn" d="M500 606 Q420 620 390 642" marker-end="url(#arr)"/>
<path class="conn" d="M290 482 Q300 560 330 642" marker-end="url(#arr)"/>
<!-- Legend -->
<rect x="56" y="710" width="608" height="60" rx="8" fill="none" stroke="var(--border)" stroke-width="1"/>
<text x="68" y="728" class="t-title">图例</text>
<rect x="68" y="738" width="20" height="12" rx="3" class="node-brand"/>
<text x="94" y="748" class="t-mini">核心组件 / 主路径</text>
<rect x="220" y="738" width="20" height="12" rx="3" class="node-neutral"/>
<text x="246" y="748" class="t-mini">普通组件</text>
<rect x="360" y="738" width="20" height="12" rx="3" class="node-external"/>
<text x="386" y="748" class="t-mini">外部服务</text>
<line x1="480" y1="744" x2="510" y2="744" class="conn-brand"/>
<text x="516" y="748" class="t-mini">工具调用主链路</text>
</svg>
</div>
### 智能网关路由
通过网关助手(GatewayAssistant)实现 LLM 智能路由决策:用户用自然语言提问,LLM 自动判断意图并路由到对应领域助手,全程无需手动选择入口。
@@ -50,7 +225,7 @@
| 助手 | 路由路径 | 覆盖能力 |
|------|---------|---------|
| 摸鱼小助手 | `/ai/assist/chat` | 通用问答,平台介绍 |
| 摸鱼小助手 | `/ai/agent/chat` | 通用问答,平台介绍 |
| 认证中心助手 | `/ai/oauth2/assist` | 菜单、组织、角色、帐号、租户管理 |
| 系统中心助手 | `/ai/sys/assist` | 字典、分类目录、日志、代码生成、数据库、数据源 |
| 低代码中心助手 | `/ai/nocode/assist` | 自助大屏、自助API、组件管理、公式信息 |
+24
View File
@@ -34,4 +34,28 @@ CREATE TABLE `ai_model_config` (
PRIMARY KEY (`id`) USING BTREE
) ENGINE = InnoDB CHARACTER SET = utf8mb4 COLLATE = utf8mb4_general_ci COMMENT = 'AI模型配置信息' ROW_FORMAT = Dynamic;
-- ----------------------------
-- Table structure for ai_mcp_server_config
-- ----------------------------
DROP TABLE IF EXISTS `ai_mcp_server_config`;
CREATE TABLE `ai_mcp_server_config` (
`id` varchar(36) CHARACTER SET utf8mb4 COLLATE utf8mb4_general_ci NOT NULL COMMENT '主键',
`server_name` varchar(100) CHARACTER SET utf8mb4 COLLATE utf8mb4_general_ci NOT NULL COMMENT 'MCP服务器名称',
`transport_type` varchar(20) CHARACTER SET utf8mb4 COLLATE utf8mb4_general_ci NOT NULL COMMENT '传输类型: stdio/sse/streamable',
`command` varchar(500) CHARACTER SET utf8mb4 COLLATE utf8mb4_general_ci NULL DEFAULT NULL COMMENT 'stdio模式启动命令(如node/python)',
`args` text CHARACTER SET utf8mb4 COLLATE utf8mb4_general_ci NULL COMMENT 'stdio模式参数(JSON数组,如["server.js","--port","3000"])',
`env` text CHARACTER SET utf8mb4 COLLATE utf8mb4_general_ci NULL COMMENT 'stdio模式环境变量(JSON对象)',
`sse_url` varchar(500) CHARACTER SET utf8mb4 COLLATE utf8mb4_general_ci NULL DEFAULT NULL COMMENT '远程服务基础URL(sse/streamable通用,如https://mcp.example.com)',
`sse_endpoint` varchar(200) CHARACTER SET utf8mb4 COLLATE utf8mb4_general_ci NULL DEFAULT NULL COMMENT '远程服务端点路径(sse/streamable通用,如/sse或/mcp)',
`auth_token` varchar(500) CHARACTER SET utf8mb4 COLLATE utf8mb4_general_ci NULL DEFAULT NULL COMMENT '认证Token(Bearer)',
`status` tinyint NOT NULL DEFAULT 0 COMMENT '状态 0正常 1停用',
`remark` varchar(500) CHARACTER SET utf8mb4 COLLATE utf8mb4_general_ci NULL DEFAULT NULL COMMENT '备注',
`create_by` varchar(50) CHARACTER SET utf8mb4 COLLATE utf8mb4_general_ci NULL DEFAULT NULL COMMENT '创建人',
`create_time` datetime NULL DEFAULT NULL COMMENT '创建时间',
`update_by` varchar(50) CHARACTER SET utf8mb4 COLLATE utf8mb4_general_ci NULL DEFAULT NULL COMMENT '修改人',
`update_time` datetime NULL DEFAULT NULL COMMENT '修改时间',
PRIMARY KEY (`id`) USING BTREE,
UNIQUE INDEX `server_name_index`(`server_name`) USING BTREE
) ENGINE = InnoDB CHARACTER SET = utf8mb4 COLLATE = utf8mb4_general_ci COMMENT = 'MCP服务器配置表' ROW_FORMAT = Dynamic;
SET FOREIGN_KEY_CHECKS = 1;
+24
View File
@@ -1986,4 +1986,28 @@ CREATE TABLE `ai_model_config` (
PRIMARY KEY (`id`) USING BTREE
) ENGINE = InnoDB CHARACTER SET = utf8mb4 COLLATE = utf8mb4_general_ci COMMENT = 'AI模型配置信息' ROW_FORMAT = Dynamic;
-- ----------------------------
-- Table structure for ai_mcp_server_config
-- ----------------------------
DROP TABLE IF EXISTS `ai_mcp_server_config`;
CREATE TABLE `ai_mcp_server_config` (
`id` varchar(36) CHARACTER SET utf8mb4 COLLATE utf8mb4_general_ci NOT NULL COMMENT '主键',
`server_name` varchar(100) CHARACTER SET utf8mb4 COLLATE utf8mb4_general_ci NOT NULL COMMENT 'MCP服务器名称',
`transport_type` varchar(20) CHARACTER SET utf8mb4 COLLATE utf8mb4_general_ci NOT NULL COMMENT '传输类型: stdio/sse/streamable',
`command` varchar(500) CHARACTER SET utf8mb4 COLLATE utf8mb4_general_ci NULL DEFAULT NULL COMMENT 'stdio模式启动命令(如node/python)',
`args` text CHARACTER SET utf8mb4 COLLATE utf8mb4_general_ci NULL COMMENT 'stdio模式参数(JSON数组,如["server.js","--port","3000"])',
`env` text CHARACTER SET utf8mb4 COLLATE utf8mb4_general_ci NULL COMMENT 'stdio模式环境变量(JSON对象)',
`sse_url` varchar(500) CHARACTER SET utf8mb4 COLLATE utf8mb4_general_ci NULL DEFAULT NULL COMMENT '远程服务基础URL(sse/streamable通用,如https://mcp.example.com)',
`sse_endpoint` varchar(200) CHARACTER SET utf8mb4 COLLATE utf8mb4_general_ci NULL DEFAULT NULL COMMENT '远程服务端点路径(sse/streamable通用,如/sse或/mcp)',
`auth_token` varchar(500) CHARACTER SET utf8mb4 COLLATE utf8mb4_general_ci NULL DEFAULT NULL COMMENT '认证Token(Bearer)',
`status` tinyint NOT NULL DEFAULT 0 COMMENT '状态 0正常 1停用',
`remark` varchar(500) CHARACTER SET utf8mb4 COLLATE utf8mb4_general_ci NULL DEFAULT NULL COMMENT '备注',
`create_by` varchar(50) CHARACTER SET utf8mb4 COLLATE utf8mb4_general_ci NULL DEFAULT NULL COMMENT '创建人',
`create_time` datetime NULL DEFAULT NULL COMMENT '创建时间',
`update_by` varchar(50) CHARACTER SET utf8mb4 COLLATE utf8mb4_general_ci NULL DEFAULT NULL COMMENT '修改人',
`update_time` datetime NULL DEFAULT NULL COMMENT '修改时间',
PRIMARY KEY (`id`) USING BTREE,
UNIQUE INDEX `server_name_index`(`server_name`) USING BTREE
) ENGINE = InnoDB CHARACTER SET = utf8mb4 COLLATE = utf8mb4_general_ci COMMENT = 'MCP服务器配置表' ROW_FORMAT = Dynamic;
SET FOREIGN_KEY_CHECKS = 1;
@@ -0,0 +1,59 @@
package cn.com.mfish.ai.api.entity;
import cn.com.mfish.common.core.entity.BaseEntity;
import cn.idev.excel.annotation.ExcelProperty;
import com.baomidou.mybatisplus.annotation.IdType;
import com.baomidou.mybatisplus.annotation.TableId;
import com.baomidou.mybatisplus.annotation.TableName;
import io.swagger.v3.oas.annotations.media.Schema;
import lombok.Data;
import lombok.EqualsAndHashCode;
import lombok.experimental.Accessors;
/**
* @description: MCP服务器配置信息
* @author: mfish
* @date: 2026-07-21
* @version: V2.4.1
*/
@Data
@TableName("ai_mcp_server_config")
@EqualsAndHashCode(callSuper = true)
@Schema(description = "ai_mcp_server_config对象 MCP服务器配置信息")
public class McpServerConfig extends BaseEntity<String> {
@ExcelProperty("唯一ID")
@Schema(description = "唯一ID")
@TableId(type = IdType.ASSIGN_UUID)
@Accessors(chain = true)
private String id;
@ExcelProperty("MCP服务器名称")
@Schema(description = "MCP服务器名称")
private String serverName;
@ExcelProperty("传输类型: stdio/sse/streamable")
@Schema(description = "传输类型: stdio/sse/streamable")
private String transportType;
@ExcelProperty("stdio模式启动命令")
@Schema(description = "stdio模式启动命令(如node/python)")
private String command;
@ExcelProperty("stdio模式参数(JSON数组)")
@Schema(description = "stdio模式参数(JSON数组,如[\"server.js\",\"--port\",\"3000\"])")
private String args;
@ExcelProperty("stdio模式环境变量(JSON对象)")
@Schema(description = "stdio模式环境变量(JSON对象)")
private String env;
@ExcelProperty("远程服务基础URL(sse/streamable通用)")
@Schema(description = "远程服务基础URL(sse/streamable通用,如https://mcp.example.com)")
private String sseUrl;
@ExcelProperty("远程服务端点路径(sse/streamable通用)")
@Schema(description = "远程服务端点路径(sse/streamable通用,如/sse或/mcp)")
private String sseEndpoint;
@ExcelProperty("认证Token")
@Schema(description = "认证Token(Bearer)")
private String authToken;
@ExcelProperty("状态 0正常 1停用")
@Schema(description = "状态 0正常 1停用")
private Short status;
@ExcelProperty("备注")
@Schema(description = "备注")
private String remark;
}
@@ -9,9 +9,11 @@ import cn.com.mfish.common.ai.entity.AiRequest;
import cn.com.mfish.common.ai.entity.ChatResponseVo;
import cn.com.mfish.common.ai.entity.EventType;
import cn.com.mfish.common.ai.entity.PlanStep;
import cn.com.mfish.common.ai.memory.ConversationMemory;
import cn.com.mfish.common.ai.memory.ConversationMemoryStore;
import cn.com.mfish.common.ai.memory.DocumentChunk;
import cn.com.mfish.common.core.utils.AuthInfoUtils;
import cn.com.mfish.common.core.utils.ServletUtils;
import cn.com.mfish.common.core.utils.StringUtils;
import lombok.extern.slf4j.Slf4j;
import org.springframework.stereotype.Component;
import org.springframework.web.context.request.RequestAttributes;
@@ -51,30 +53,34 @@ public class AgentRuntime {
private final Executor executor;
private final LlmModelRouter llmModelRouter;
private final FileParseService fileParseService;
private final ConversationMemoryStore memoryStore;
public AgentRuntime(Planner planner, Executor executor, LlmModelRouter llmModelRouter,
FileParseService fileParseService) {
FileParseService fileParseService, ConversationMemoryStore memoryStore) {
this.planner = planner;
this.executor = executor;
this.llmModelRouter = llmModelRouter;
this.fileParseService = fileParseService;
this.memoryStore = memoryStore;
}
/**
* 运行智能体编排
* <p>
* 在请求线程捕获租户上下文快照,后台启动 Plan → Execute 流程,同时返回事件流供前端订阅。
* 编排流程在 boundedElastic 调度器上执行,不阻塞当前线程。
* 重构后流程(接入 Memory 模块):
* <ol>
* <li>获取(或创建)sessionId 对应的 {@link ConversationMemory}</li>
* <li>在请求线程捕获 TenantContext 并绑定到 Memory 的 Vars Context</li>
* <li>若请求携带 fileIds:在请求线程解析文件为 DocumentChunk 列表,
* 注入到 Memory 的 Document Context</li>
* <li>从 Memory 读取 Document Context 拼接文本,注入到用户 prompt 前</li>
* <li>后台启动 Plan → Execute 流程</li>
* </ol>
* </p>
* <p>
* 关键点:必须在请求线程捕获 {@link TenantContext}(含 tenantId/userId/token/RequestAttributes/Exchange),
* 因为 AuthInfoUtils 和 RequestContextHolder 不支持异步线程访问,编排切到 boundedElastic 后
* 需要用快照构建 ToolContext 和 ChatClient
* </p>
* <p>
* 文件解析:若请求携带 fileIds,在请求线程同步加载文件内容(Feign 调用 mf-storage),
* 拼接到用户提示词前。必须在请求线程执行,因为 Feign 的 BearerTokenInterceptor 依赖
* RequestContextHolder 中继令牌,异步线程拿不到。
* 关键点:必须在请求线程完成 Memory 写入(TenantContext 绑定 + 文件解析),
* 因为 AuthInfoUtils 和 Feign BearerTokenInterceptor 都依赖 RequestContextHolder
* 切到 boundedElastic 异步线程后拿不到请求上下文
* </p>
* <p>
* 所有返回的 {@link ChatResponseVo} 的 id 字段填充为 {@link AiRequest#getId()}
@@ -88,14 +94,25 @@ public class AgentRuntime {
String requestId = aiRequest.getId();
String sessionId = aiRequest.getSessionId();
String prompt = aiRequest.getMessage() != null ? aiRequest.getMessage().getContent() : null;
// 用请求 id 作为事件 id,与普通聊天返回结构一致
EventBus eventBus = new EventBus(requestId);
// 在请求线程捕获租户上下文快照,供异步编排使用
TenantContext tenantContext = captureTenantContext();
// 获取(或创建)会话 Memory
ConversationMemory memory = memoryStore.getOrCreate(sessionId);
// 文件解析必须在请求线程执行:Feign BearerTokenInterceptor 依赖 RequestContextHolder
prompt = resolveFileContents(prompt, aiRequest.getFileIds());
// 在请求线程捕获租户上下文并绑定到 Memory 的 Vars Context
TenantContext tenantContext = captureTenantContext();
memory.bindTenantContext(tenantContext);
// 文件解析 + 注入 Memory 的 Document Context(必须在请求线程执行)
List<DocumentChunk> chunks = fileParseService.loadAsChunks(aiRequest.getFileIds());
if (!chunks.isEmpty()) {
memory.addDocumentChunks(chunks);
}
// Document Context 的拼接交由 Planner 统一处理:
// Planner.plan() 会从 Memory 读取 getSystemContext()(含 Vars + Document)拼接到 prompt 前。
// AgentRuntime 只负责 Memory 写入(文件解析 + 租户绑定),不再拼接 DocumentContext 到 prompt
// 避免与 Planner 重复拼接。
// 后台启动编排流程
runOrchestration(sessionId, prompt, eventBus, tenantContext);
@@ -104,31 +121,6 @@ public class AgentRuntime {
return eventBus.asFlux();
}
/**
* 加载文件内容并拼接到提示词前
* <p>
* 若 fileIds 为空或全部加载失败,返回原始 prompt。
* 文件内容加载失败时记录日志但不阻断流程,降级为纯文本对话。
* </p>
*
* @param prompt 原始用户提示词
* @param fileIds 文件fileKey列表
* @return 拼接后的提示词
*/
private String resolveFileContents(String prompt, List<String> fileIds) {
if (fileIds == null || fileIds.isEmpty()) {
return prompt;
}
String fileContents = fileParseService.loadFileContents(fileIds);
if (StringUtils.isEmpty(fileContents)) {
log.warn("[AgentRuntime] 文件内容加载为空 fileIds={}", fileIds);
return prompt;
}
return "以下是用户上传的文件内容,请基于文件内容进行分析:\n\n"
+ fileContents
+ "\n用户需求:" + prompt;
}
/**
* 在请求线程捕获租户上下文快照
* <p>
@@ -7,6 +7,8 @@ import cn.com.mfish.common.ai.engine.ApiToolEngine;
import cn.com.mfish.common.ai.entity.AiRequest;
import cn.com.mfish.common.ai.entity.ChatResponseVo;
import cn.com.mfish.common.ai.agent.ToolCapable;
import cn.com.mfish.common.ai.memory.ConversationMemory;
import cn.com.mfish.common.ai.memory.ConversationMemoryStore;
import cn.com.mfish.common.core.constants.RPCConstants;
import cn.com.mfish.common.core.utils.AuthInfoUtils;
import cn.com.mfish.common.core.utils.ServletUtils;
@@ -61,6 +63,19 @@ public abstract class BaseAssistant implements IClientAssistant, ToolCapable {
@Autowired
protected FileParseService fileParseService;
/**
* 会话记忆存储:三驾马车 Memory 模块的入口
* <p>
* 用于按 sessionId 获取 {@link ConversationMemory},将文件解析结果、租户上下文、
* 业务变量等统一收纳到 Memory,供后续 Planner 拼接上下文。
* </p>
* <p>
* 字段注入避免修改所有子类构造函数。
* </p>
*/
@Autowired
protected ConversationMemoryStore memoryStore;
public BaseAssistant(ChatMemory chatMemory, LlmModelRouter llmModelRouter, ApiToolEngine apiToolEngine) {
this.llmModelRouter = llmModelRouter;
this.chatMemory = chatMemory;
@@ -279,25 +294,47 @@ public abstract class BaseAssistant implements IClientAssistant, ToolCapable {
/**
* 聊天返回id
* <p>
* 若请求携带fileIds,会先通过FileParseService从文件服务获取文件内容,
* 将其拼接到用户提示词前,再交由子类chat(sessionId, prompt)处理。
* 重构后流程(接入 Memory 模块):
* <ol>
* <li>获取(或创建)sessionId 对应的 {@link ConversationMemory}</li>
* <li>若请求携带 fileIds:通过 FileParseService 解析为 DocumentChunk 列表,
* 注入到 Memory 的 Document Context(替代旧的字符串拼接)</li>
* <li>从 Memory 读取 Document Context 拼接文本,注入到用户 prompt 前</li>
* <li>交由子类 chat(sessionId, prompt) 处理(保持向后兼容)</li>
* </ol>
* <p>
* 关键约束:文件解析必须在请求线程执行(Feign BearerTokenInterceptor 依赖
* RequestContextHolder),Memory 写入也在请求线程完成(保证一致性)。
*
* @return 聊天信息
*/
@Override
public Flux<ChatResponseVo> chat(AiRequest aiRequest) {
String sessionId = aiRequest.getSessionId();
String prompt = aiRequest.getMessage().getContent();
ConversationMemory memory = memoryStore.getOrCreate(sessionId);
// 文件解析 + 注入 Memory 的 Document Context
List<String> fileIds = aiRequest.getFileIds();
if (fileIds != null && !fileIds.isEmpty()) {
String fileContents = fileParseService.loadFileContents(fileIds);
if (StringUtils.isNotEmpty(fileContents)) {
prompt = "以下是用户上传的文件内容,请基于文件内容进行分析:\n\n"
+ fileContents
+ "\n用户问题:" + prompt;
List<cn.com.mfish.common.ai.memory.DocumentChunk> chunks = fileParseService.loadAsChunks(fileIds);
if (!chunks.isEmpty()) {
memory.addDocumentChunks(chunks);
}
}
// 从 Memory 读取文档上下文,拼接到 prompt 前
// 当前使用 getDocumentContext() 全量注入;当文档较多、token 预算紧张时,
// 未来可切换为 memory.searchDocumentContext(prompt, 5) 走 RAG 检索(向量库版 Memory 覆写此方法)。
String documentContext = memory.getDocumentContext();
if (StringUtils.isNotEmpty(documentContext)) {
prompt = "以下是用户上传的文件内容,请基于文件内容进行分析:\n\n"
+ documentContext
+ "\n用户问题:" + prompt;
}
final String finalPrompt = prompt;
return chat(aiRequest.getSessionId(), finalPrompt)
return chat(sessionId, finalPrompt)
.filter(resp -> "STOP".equals(Objects.requireNonNull(resp.getResult()).getMetadata().getFinishReason())
|| StringUtils.isNotEmpty(resp.getResult().getOutput().getText()))
.map(resp -> new ChatResponseVo().setId(aiRequest.getId())
@@ -2,9 +2,13 @@ package cn.com.mfish.ai.agent;
import cn.com.mfish.ai.service.LlmModelRouter;
import cn.com.mfish.common.ai.agent.TenantContext;
import cn.com.mfish.common.ai.capability.ActionDefinition;
import cn.com.mfish.common.ai.capability.CapabilityEngine;
import cn.com.mfish.common.ai.entity.AgentPlan;
import cn.com.mfish.common.ai.entity.PlanStep;
import cn.com.mfish.common.core.constants.ServiceConstants;
import cn.com.mfish.common.ai.memory.ConversationMemory;
import cn.com.mfish.common.ai.memory.ConversationMemoryStore;
import cn.com.mfish.common.core.utils.StringUtils;
import lombok.extern.slf4j.Slf4j;
import org.springframework.ai.chat.client.ChatClient;
import org.springframework.ai.chat.client.advisor.MessageChatMemoryAdvisor;
@@ -14,24 +18,37 @@ import org.springframework.stereotype.Component;
import reactor.core.publisher.Mono;
import reactor.core.scheduler.Schedulers;
import java.util.Arrays;
import java.util.LinkedHashMap;
import java.util.List;
import java.util.Map;
import java.util.Objects;
import java.util.Set;
import java.util.stream.Collectors;
import static org.springframework.ai.chat.memory.ChatMemory.CONVERSATION_ID;
/**
* 任务规划器
* 任务规划器(已接入 CapabilityEngine + Memory
* <p>
* 调用 LLM 将用户原始需求拆解为结构化的 {@link AgentPlan}
* 重构后职责(第二阶段存量平移)
* <ol>
* <li>分析用户意图</li>
* <li>拆解为多个可执行的步骤</li>
* <li>每步指定需要调用的微服务集合</li>
* <li>从 {@link ConversationMemoryStore} 读取会话 Memory,获取系统上下文(Vars + Document</li>
* <li>从 {@link CapabilityEngine} 获取所有可用 {@link ActionDefinition},按服务分组展示给 LLM</li>
* <li>将上下文 + 动作列表 + 用户需求打包发给 LLM,生成结构化 {@link AgentPlan}</li>
* </ol>
* </p>
* <p>
* <b>与旧版差异</b>
* <ul>
* <li>服务列表来源:旧版从 {@code ServiceConstants.MfService.values()} 硬编码枚举获取;
* 新版从 {@link CapabilityEngine#getAvailableActions()} 动态获取,自动反映已注册的子引擎</li>
* <li>上下文注入:旧版仅传入原始 prompt;新版从 Memory 读取 {@code getSystemContext()}
* 拼接到 prompt 前,包含租户信息、业务变量和文档内容</li>
* <li>动作摘要:新版在规划提示词中展示每个服务的动作数量和示例动作名,
* 帮助 LLM 更精准地选择 serviceIds</li>
* </ul>
* </p>
* <p>
* 采用结构化输出(responseEntity),LLM 直接返回 JSON 反序列化为 AgentPlan。
* 必须在请求线程解析租户并构建 ChatClientAuthInfoUtils 不支持异步)。
* </p>
@@ -45,22 +62,36 @@ public class Planner {
private final ChatMemory chatMemory;
private final LlmModelRouter llmModelRouter;
private final CapabilityEngine capabilityEngine;
private final ConversationMemoryStore memoryStore;
public Planner(ChatMemory chatMemory, LlmModelRouter llmModelRouter) {
public Planner(ChatMemory chatMemory, LlmModelRouter llmModelRouter,
CapabilityEngine capabilityEngine, ConversationMemoryStore memoryStore) {
this.chatMemory = chatMemory;
this.llmModelRouter = llmModelRouter;
this.capabilityEngine = capabilityEngine;
this.memoryStore = memoryStore;
}
/**
* 规划:将用户需求拆解为步骤列表
* <p>
* 注意:本方法会被异步调度器(boundedElastic)调用,因此 ChatClient 必须在请求线程
* 预构建并通过闭包传入。{@link TenantContext} 中包含请求线程捕获的 tenantId,
* 重构后流程:
* <ol>
* <li>从 Memory 读取系统上下文(Vars + Document),拼接到 prompt 前</li>
* <li>从 CapabilityEngine 获取所有可用 ActionDefinition,构建规划提示词</li>
* <li>调用 LLM 生成结构化 AgentPlan</li>
* <li>失败时降级为单步执行,聚合所有已注册服务</li>
* </ol>
* </p>
* <p>
* 注意:本方法会被异步调度器(boundedElastic)调用,因此 ChatClient 和系统提示词
* 必须在请求线程预构建并通过闭包传入。{@link TenantContext} 中包含请求线程捕获的 tenantId,
* 用于路由到该租户的 ChatModel。
* </p>
*
* @param sessionId 会话ID
* @param prompt 用户原始需求
* @param prompt 用户原始需求(已由 AgentRuntime 拼接 Document Context
* @param tenantContext 请求线程捕获的租户上下文
* @return 执行计划
*/
@@ -70,16 +101,26 @@ public class Planner {
? tenantContext.getTenantId()
: llmModelRouter.currentTenantId();
ChatClient chatClient = getChatClient(tenantId);
// 从 Memory 读取系统上下文(Vars + Document),拼接到 prompt 前
String systemContext = resolveSystemContext(sessionId);
// 从 CapabilityEngine 获取动作列表,构建规划提示词
String systemPrompt = buildPlannerPrompt();
// 拼接最终 prompt:系统上下文 + 用户需求
String finalPrompt = StringUtils.isNotEmpty(systemContext)
? systemContext + "\n用户需求:" + prompt
: prompt;
return Mono.fromCallable(() -> {
var responseEntity = chatClient.prompt()
.system(systemPrompt)
.user(prompt)
.user(finalPrompt)
.advisors(a -> a.param(CONVERSATION_ID, sessionId))
.call()
.responseEntity(AgentPlan.class);
AgentPlan plan = Objects.requireNonNullElseGet(responseEntity.entity(), () -> fallbackPlan(prompt));
AgentPlan plan = Objects.requireNonNullElseGet(responseEntity.entity(), () -> fallbackPlan(finalPrompt));
plan.setOriginalPrompt(prompt);
log.info("[Planner] 规划完成, 步骤数={}, summary={}",
plan.getSteps() != null ? plan.getSteps().size() : 0, plan.getSummary());
@@ -88,17 +129,63 @@ public class Planner {
.subscribeOn(Schedulers.boundedElastic())
.onErrorResume(ex -> {
log.error("[Planner] 规划失败,降级为单步执行", ex);
return Mono.just(fallbackPlan(prompt));
return Mono.just(fallbackPlan(finalPrompt));
});
}
/**
* 构建规划师系统提示词,包含可用服务列表
* 从 Memory 读取系统上下文(Vars Context + Document Context
* <p>
* AgentRuntime 已在请求线程将 DocumentContext 拼接到 prompt 中,
* 此处进一步补充 Vars Context(租户信息、业务变量),
* 使 Planner 能感知用户身份和业务上下文。
* </p>
*/
private String resolveSystemContext(String sessionId) {
try {
ConversationMemory memory = memoryStore.getOrCreate(sessionId);
return memory.getSystemContext();
} catch (Exception e) {
log.warn("[Planner] 读取 Memory 系统上下文失败 sessionId={}", sessionId, e);
return "";
}
}
/**
* 构建规划师系统提示词,包含可用服务列表和动作摘要
* <p>
* 重构后从 {@link CapabilityEngine#getAvailableActions()} 动态获取动作列表,
* 按服务分组展示。相比旧版从 {@code ServiceConstants.MfService.values()} 硬编码枚举获取,
* 新版自动反映已注册的子引擎状态(如某服务未启动则不出现在列表中)。
* </p>
*/
private String buildPlannerPrompt() {
String serviceList = Arrays.stream(ServiceConstants.MfService.values())
.map(s -> String.format(" - %s: %s", s.getValue(), s.getGatewayPrefix()))
.collect(Collectors.joining("\n"));
// 从 CapabilityEngine 获取所有动作,按 serviceId 分组
List<ActionDefinition> actions = capabilityEngine.getAvailableActions();
Map<String, List<ActionDefinition>> actionsByService = groupActionsByService(actions);
// 构建服务列表(仅包含有动作的服务)
StringBuilder serviceList = new StringBuilder();
for (Map.Entry<String, List<ActionDefinition>> entry : actionsByService.entrySet()) {
String serviceId = entry.getKey();
List<ActionDefinition> serviceActions = entry.getValue();
serviceList.append(String.format(" - %s: %d 个可用动作", serviceId, serviceActions.size()));
// 展示前 3 个动作名作为示例,帮助 LLM 理解服务能力
List<String> sampleNames = serviceActions.stream()
.limit(3)
.map(ActionDefinition::getName)
.toList();
if (!sampleNames.isEmpty()) {
serviceList.append("(示例: ").append(String.join(", ", sampleNames)).append("");
}
serviceList.append("\n");
}
// 如果没有动作,使用兜底提示
if (serviceList.isEmpty()) {
serviceList.append(" (当前无可用服务,请直接回答用户问题)\n");
}
return """
你是"摸鱼低代码"平台的任务规划师。
@@ -154,17 +241,39 @@ public class Planner {
}
/**
* 规划失败时的兜底:单步执行原始需求,聚合所有服务工具
* 将动作列表按 serviceId 分组(保持注册顺序)
*/
private Map<String, List<ActionDefinition>> groupActionsByService(List<ActionDefinition> actions) {
Map<String, List<ActionDefinition>> grouped = new LinkedHashMap<>();
if (actions == null || actions.isEmpty()) {
return grouped;
}
for (ActionDefinition action : actions) {
String serviceId = action.getServiceId() != null ? action.getServiceId() : "unknown";
grouped.computeIfAbsent(serviceId, k -> new java.util.ArrayList<>()).add(action);
}
return grouped;
}
/**
* 规划失败时的兜底:单步执行原始需求,聚合所有已注册服务
* <p>
* 重构后从 {@link CapabilityEngine#getAvailableActions()} 提取所有 serviceId
* 相比旧版从 {@code ServiceConstants.MfService.values()} 硬编码枚举获取,
* 新版仅聚合实际有动作的服务,避免向未启动的服务发送请求。
* </p>
*/
private AgentPlan fallbackPlan(String prompt) {
List<String> allServices = Arrays.stream(ServiceConstants.MfService.values())
.map(ServiceConstants.MfService::getValue)
.collect(Collectors.toList());
Set<String> allServiceIds = capabilityEngine.getAvailableActions().stream()
.map(ActionDefinition::getServiceId)
.filter(Objects::nonNull)
.collect(Collectors.toSet());
List<String> serviceIdList = List.copyOf(allServiceIds);
return new AgentPlan()
.setOriginalPrompt(prompt)
.setSummary("规划降级:直接执行")
.setSteps(List.of(
new PlanStep(prompt, allServices)
new PlanStep(prompt, serviceIdList)
));
}
@@ -0,0 +1,69 @@
package cn.com.mfish.ai.controller;
import cn.com.mfish.ai.api.entity.McpServerConfig;
import cn.com.mfish.ai.service.McpServerConfigService;
import cn.com.mfish.common.core.web.PageResult;
import cn.com.mfish.common.core.web.ReqPage;
import cn.com.mfish.common.core.web.Result;
import cn.com.mfish.common.core.enums.OperateType;
import cn.com.mfish.common.log.annotation.Log;
import cn.com.mfish.common.oauth.annotation.RequiresPermissions;
import io.swagger.v3.oas.annotations.Operation;
import io.swagger.v3.oas.annotations.tags.Tag;
import jakarta.annotation.Resource;
import lombok.extern.slf4j.Slf4j;
import org.springframework.web.bind.annotation.*;
/**
* @description: MCP服务器配置信息
* @author: mfish
* @date: 2026-07-21
* @version: V2.4.1
*/
@Tag(name = "MCP服务器配置")
@RestController
@RequestMapping("/mcp/config")
@Slf4j
public class McpServerConfigController {
@Resource
McpServerConfigService mcpServerConfigService;
@Operation(summary = "MCP服务器配置-分页列表查询", description = "MCP服务器配置-分页列表查询")
@GetMapping
@RequiresPermissions("ai:mcp:query")
public Result<PageResult<McpServerConfig>> queryPageList(McpServerConfig req, ReqPage reqPage) {
return mcpServerConfigService.queryPageList(req, reqPage);
}
@Log(title = "MCP服务器配置-新增", operateType = OperateType.INSERT)
@Operation(summary = "MCP服务器配置-新增", description = "MCP服务器配置-新增")
@PostMapping
@RequiresPermissions("ai:mcp:insert")
public Result<McpServerConfig> add(@RequestBody McpServerConfig entity) {
return mcpServerConfigService.insert(entity);
}
@Log(title = "MCP服务器配置-修改", operateType = OperateType.UPDATE)
@Operation(summary = "MCP服务器配置-修改", description = "MCP服务器配置-修改")
@PutMapping
@RequiresPermissions("ai:mcp:update")
public Result<McpServerConfig> edit(@RequestBody McpServerConfig entity) {
return mcpServerConfigService.update(entity);
}
@Log(title = "MCP服务器配置-通过id删除", operateType = OperateType.DELETE)
@Operation(summary = "MCP服务器配置-通过id删除", description = "MCP服务器配置-通过id删除")
@DeleteMapping("/{id}")
@RequiresPermissions("ai:mcp:delete")
public Result<Boolean> delete(@PathVariable String id) {
return mcpServerConfigService.delete(id);
}
@Operation(summary = "MCP服务器配置-通过id查询", description = "MCP服务器配置-通过id查询")
@GetMapping("/{id}")
@RequiresPermissions("ai:mcp:query")
public Result<McpServerConfig> getById(@PathVariable String id) {
return Result.ok(mcpServerConfigService.getById(id), "查询成功");
}
}
@@ -0,0 +1,14 @@
package cn.com.mfish.ai.mapper;
import cn.com.mfish.ai.api.entity.McpServerConfig;
import com.baomidou.mybatisplus.core.mapper.BaseMapper;
/**
* @description: MCP服务器配置信息
* @author: mfish
* @date: 2026-07-21
* @version: V2.4.1
*/
public interface McpServerConfigMapper extends BaseMapper<McpServerConfig> {
}
@@ -1,5 +1,6 @@
package cn.com.mfish.ai.service;
import cn.com.mfish.common.ai.memory.DocumentChunk;
import cn.com.mfish.common.core.constants.RPCConstants;
import cn.com.mfish.common.core.utils.StringUtils;
import cn.com.mfish.common.storage.api.entity.StorageInfo;
@@ -14,9 +15,11 @@ import org.springframework.stereotype.Service;
import java.io.InputStream;
import java.nio.charset.StandardCharsets;
import java.util.ArrayList;
import java.util.List;
import java.util.Locale;
import java.util.Set;
import java.util.UUID;
/**
* 文件解析服务
@@ -111,6 +114,9 @@ public class FileParseService {
* <p>
* 遍历fileKey列表,逐个获取文件元数据与内容,将文本/Office文档内容拼接到StringBuilder。
* 单个文件读取失败不影响其他文件,异常会被记录日志并跳过。
* <p>
* 旧版接口,保留向后兼容。新调用方建议使用 {@link #loadAsChunks} 获取结构化的
* {@link DocumentChunk} 列表,再注入到 Memory 模块。
*
* @param fileIds 文件fileKey列表
* @return 拼接好的文件内容提示词片段;列表为空或全部失败时返回空字符串
@@ -141,6 +147,45 @@ public class FileParseService {
return sb.toString();
}
/**
* 加载文件并解析为 {@link DocumentChunk} 列表
* <p>
* 与 {@link #loadFileContents} 的区别:返回结构化的文档块列表,便于 Memory 模块
* 按 chunk 维度管理文档上下文,支持后续按相关性检索(RAG)、按 token 上限分片注入等扩展。
* </p>
* <p>
* 当前实现:每个文件产出 1 个 Chunk(小文件场景),未来可扩展为按段落/页码/字符数分片。
* Chunk 内容已截断到 {@link #MAX_CONTENT_CHARS},可直接注入 LLM 上下文。
* </p>
* <p>
* 必须在请求线程调用:内部通过 Feign 调用 mf-storageBearerTokenInterceptor 依赖
* RequestContextHolder 中继令牌。
* </p>
*
* @param fileIds 文件fileKey列表
* @return 文档块列表;列表为空或全部失败时返回空列表(不会返回 null)
*/
public List<DocumentChunk> loadAsChunks(List<String> fileIds) {
if (fileIds == null || fileIds.isEmpty()) {
return List.of();
}
List<DocumentChunk> chunks = new ArrayList<>(fileIds.size());
for (String fileKey : fileIds) {
if (StringUtils.isEmpty(fileKey)) {
continue;
}
try {
DocumentChunk chunk = loadOneFileAsChunk(fileKey);
if (chunk != null && StringUtils.isNotEmpty(chunk.getContent())) {
chunks.add(chunk);
}
} catch (Exception e) {
log.warn("加载文件块失败 fileKey={} reason={}", fileKey, e.getMessage());
}
}
return chunks;
}
/**
* 加载单个文件并构建内容片段
* <p>
@@ -176,6 +221,47 @@ public class FileParseService {
return buildBinaryFileHint(fileName, fileType, storageInfo.getFileSize());
}
/**
* 加载单个文件并构建为 {@link DocumentChunk}
* <p>
* 与 {@link #loadOneFile} 共享文件元数据查询和解析路径,但返回结构化 Chunk 而非拼接文本。
* 二进制文件(无法解析内容的)返回 null 而非提示片段,避免无效 Chunk 进入 Memory。
* </p>
*
* @param fileKey 文件key
* @return 文档块;文件不存在或无法解析返回 null
*/
private DocumentChunk loadOneFileAsChunk(String fileKey) {
Result<StorageInfo> infoResult = remoteStorageService.queryByKey(RPCConstants.INNER, fileKey);
if (infoResult == null || !infoResult.isSuccess() || infoResult.getData() == null) {
log.warn("文件信息查询失败 fileKey={} msg={}", fileKey,
infoResult == null ? "result is null" : infoResult.getMsg());
return null;
}
StorageInfo storageInfo = infoResult.getData();
String fileName = StringUtils.isEmpty(storageInfo.getFileName()) ? fileKey : storageInfo.getFileName();
String fileType = storageInfo.getFileType();
String content = null;
if (isTextFile(fileType, fileName)) {
content = readFileText(fileKey);
} else if (isOfficeDocument(fileType, fileName)) {
content = extractOfficeText(fileKey, fileName);
}
if (content == null || content.isEmpty()) {
return null;
}
String truncated = truncate(content);
return new DocumentChunk()
.setChunkId(UUID.randomUUID().toString())
.setFileKey(fileKey)
.setFileName(fileName)
.setFileType(fileType)
.setContent(truncated)
.setChunkIndex(0)
.setLength(truncated.length());
}
/**
* 通过Feign获取文件资源并读取为UTF-8文本
*
@@ -0,0 +1,40 @@
package cn.com.mfish.ai.service;
import cn.com.mfish.ai.api.entity.McpServerConfig;
import cn.com.mfish.common.ai.capability.McpServerInfo;
import cn.com.mfish.common.ai.capability.McpServerConfigProvider;
import cn.com.mfish.common.core.web.PageResult;
import cn.com.mfish.common.core.web.ReqPage;
import cn.com.mfish.common.core.web.Result;
import com.baomidou.mybatisplus.extension.service.IService;
import java.util.List;
/**
* @description: MCP服务器配置信息
* @author: mfish
* @date: 2026-07-21
* @version: V2.4.1
*/
public interface McpServerConfigService extends IService<McpServerConfig>, McpServerConfigProvider {
/**
* 分页查询
*/
Result<PageResult<McpServerConfig>> queryPageList(McpServerConfig req, ReqPage reqPage);
/**
* 新增
*/
Result<McpServerConfig> insert(McpServerConfig entity);
/**
* 修改
*/
Result<McpServerConfig> update(McpServerConfig entity);
/**
* 删除
*/
Result<Boolean> delete(String id);
}
@@ -0,0 +1,97 @@
package cn.com.mfish.ai.service.impl;
import cn.com.mfish.ai.api.entity.McpServerConfig;
import cn.com.mfish.ai.mapper.McpServerConfigMapper;
import cn.com.mfish.ai.service.McpServerConfigService;
import cn.com.mfish.common.ai.capability.McpServerInfo;
import cn.com.mfish.common.core.utils.StringUtils;
import cn.com.mfish.common.core.utils.Utils;
import cn.com.mfish.common.core.web.PageResult;
import cn.com.mfish.common.core.web.ReqPage;
import cn.com.mfish.common.core.web.Result;
import com.baomidou.mybatisplus.core.conditions.query.LambdaQueryWrapper;
import com.baomidou.mybatisplus.extension.service.impl.ServiceImpl;
import com.github.pagehelper.PageHelper;
import lombok.extern.slf4j.Slf4j;
import org.springframework.stereotype.Service;
import org.springframework.transaction.annotation.Transactional;
import java.util.ArrayList;
import java.util.Date;
import java.util.List;
/**
* @description: MCP服务器配置信息实现
* @author: mfish
* @date: 2026-07-21
* @version: V2.4.1
*/
@Slf4j
@Service
public class McpServerConfigServiceImpl extends ServiceImpl<McpServerConfigMapper, McpServerConfig>
implements McpServerConfigService {
@Override
public Result<PageResult<McpServerConfig>> queryPageList(McpServerConfig req, ReqPage reqPage) {
PageHelper.startPage(reqPage.getPageNum(), reqPage.getPageSize());
LambdaQueryWrapper<McpServerConfig> wrapper = new LambdaQueryWrapper<>();
if (req != null) {
wrapper.like(StringUtils.isNotEmpty(req.getServerName()), McpServerConfig::getServerName, req.getServerName());
wrapper.eq(StringUtils.isNotEmpty(req.getTransportType()), McpServerConfig::getTransportType, req.getTransportType());
wrapper.eq(req.getStatus() != null, McpServerConfig::getStatus, req.getStatus());
}
wrapper.orderByDesc(McpServerConfig::getCreateTime);
List<McpServerConfig> list = baseMapper.selectList(wrapper);
return Result.ok(new PageResult<>(list), "查询成功");
}
@Override
@Transactional
public Result<McpServerConfig> insert(McpServerConfig entity) {
if (StringUtils.isEmpty(entity.getId())) {
entity.setId(Utils.uuid32());
}
entity.setCreateTime(new Date());
baseMapper.insert(entity);
return Result.ok(entity, "新增成功");
}
@Override
@Transactional
public Result<McpServerConfig> update(McpServerConfig entity) {
entity.setUpdateTime(new Date());
baseMapper.updateById(entity);
return Result.ok(entity, "修改成功");
}
@Override
@Transactional
public Result<Boolean> delete(String id) {
baseMapper.deleteById(id);
return Result.ok(true, "删除成功");
}
/**
* 实现 {@link McpServerConfigProvider}:查询状态为"正常"的 MCP 服务器配置,
* 转换为 {@link McpServerInfo} 列表供 {@code McpCapabilityEngine} 使用
*/
@Override
public List<McpServerInfo> getActiveServerConfigs() {
LambdaQueryWrapper<McpServerConfig> wrapper = new LambdaQueryWrapper<>();
wrapper.eq(McpServerConfig::getStatus, (short) 0);
List<McpServerConfig> configs = baseMapper.selectList(wrapper);
List<McpServerInfo> result = new ArrayList<>();
for (McpServerConfig config : configs) {
result.add(new McpServerInfo()
.setServerName(config.getServerName())
.setTransportType(config.getTransportType())
.setCommand(config.getCommand())
.setArgs(config.getArgs())
.setEnv(config.getEnv())
.setSseUrl(config.getSseUrl())
.setSseEndpoint(config.getSseEndpoint())
.setAuthToken(config.getAuthToken()));
}
return result;
}
}
+11
View File
@@ -67,5 +67,16 @@
<groupId>com.baomidou</groupId>
<artifactId>mybatis-plus-spring</artifactId>
</dependency>
<!--
MCP 客户端(WebFlux Transport):McpCapabilityEngine 使用,支持 stdio + SSE 两种 Transport。
通过 Spring AI 2.0.0 的 starter 引入,版本由 spring-ai-bom 统一管理。
starter 传递依赖 mcp-core:2.0.0 + mcp-spring-webflux:2.0.0(由 Spring AI 重打包为 org.springframework.ai groupId)。
McpCapabilityEngine 手动创建 McpSyncClient(支持从数据库动态加载配置 + refresh()),
不依赖 Spring AI 的 McpClientAutoConfiguration(该配置面向 application.yml 静态配置)。
-->
<dependency>
<groupId>org.springframework.ai</groupId>
<artifactId>spring-ai-starter-mcp-client-webflux</artifactId>
</dependency>
</dependencies>
</project>
@@ -0,0 +1,66 @@
package cn.com.mfish.common.ai.capability;
import lombok.Data;
import lombok.experimental.Accessors;
/**
* 动作定义(Action Definition
* <p>
* 能力引擎向 Planner 暴露的单个可执行动作的元数据,遵循 OpenAI Tool 格式。
* 一个 {@link CapabilitySubEngine} 可暴露多个 Action{@link CapabilityEngine}
* 将所有子引擎的 Action 合并为统一列表供 Planner 选择。
* </p>
* <p>
* <b>与 Spring AI ToolCallback 的关系</b>
* <ul>
* <li>对于 TOOL 引擎:ActionDefinition 由 {@code ToolCallback.getToolDefinition()} 转换而来</li>
* <li>对于 SKILL/MCP/WORKFLOW 引擎:ActionDefinition 由各引擎自行构建</li>
* <li>Planner 可将 ActionDefinition 转换回 ToolCallback 供 ChatClient 使用,
* 也可通过 {@code CapabilityEngine.executeAction()} 直接执行</li>
* </ul>
* </p>
*
* @author: mfish
* @date: 2026/07/21
*/
@Data
@Accessors(chain = true)
public class ActionDefinition {
/**
* 动作名称(全局唯一,建议加引擎前缀避免冲突,如 {@code tool.getUserList}、{@code skill.parseFile}
*/
private String name;
/**
* 动作描述(供 LLM 理解动作用途,决定是否选择此动作)
*/
private String description;
/**
* 输入参数 JSON SchemaOpenAI Tool 格式,描述参数类型和结构)
* <p>
* 示例:
* <pre>{@code
* {
* "type": "object",
* "properties": {
* "userId": { "type": "string", "description": "用户ID" }
* },
* "required": ["userId"]
* }
* }</pre>
* </p>
*/
private String inputSchema;
/**
* 来源引擎类型
*/
private EngineType engineType;
/**
* 来源服务ID(仅 TOOL 引擎有值,标识来自哪个微服务;其他引擎为 null)
*/
private String serviceId;
}
@@ -0,0 +1,69 @@
package cn.com.mfish.common.ai.capability;
import lombok.Data;
import lombok.experimental.Accessors;
/**
* 动作执行结果
* <p>
* 由 {@link CapabilitySubEngine#execute} 返回,封装执行状态、输出内容和错误信息。
* 设计为不可变值对象:执行完成后构造,调用方只读。
* </p>
*
* @author: mfish
* @date: 2026/07/21
*/
@Data
@Accessors(chain = true)
public class ActionResult {
/**
* 是否执行成功
*/
private boolean success;
/**
* 执行输出(LLM 可读的文本,通常为 JSON 字符串或纯文本)
* <p>
* 成功时填充结果内容,失败时可为 null。
* </p>
*/
private String output;
/**
* 错误信息(失败时填充,成功时为 null)
*/
private String error;
/**
* 来源引擎类型(便于 Planner 日志追踪和结果分流处理)
*/
private EngineType engineType;
/**
* 执行耗时(毫秒)
*/
private long durationMs;
/**
* 构建成功结果
*/
public static ActionResult success(EngineType engineType, String output, long durationMs) {
return new ActionResult()
.setSuccess(true)
.setOutput(output)
.setEngineType(engineType)
.setDurationMs(durationMs);
}
/**
* 构建失败结果
*/
public static ActionResult failure(EngineType engineType, String error, long durationMs) {
return new ActionResult()
.setSuccess(false)
.setError(error)
.setEngineType(engineType)
.setDurationMs(durationMs);
}
}
@@ -0,0 +1,123 @@
package cn.com.mfish.common.ai.capability;
import cn.com.mfish.common.ai.engine.ApiToolEngine;
import lombok.extern.slf4j.Slf4j;
import org.springframework.beans.factory.SmartInitializingSingleton;
import org.springframework.boot.autoconfigure.AutoConfiguration;
import org.springframework.boot.autoconfigure.condition.ConditionalOnBean;
import org.springframework.boot.autoconfigure.condition.ConditionalOnClass;
import org.springframework.boot.autoconfigure.condition.ConditionalOnMissingBean;
import org.springframework.context.annotation.Bean;
import java.util.List;
/**
* 能力引擎自动配置
* <p>
* 注册 {@link CapabilityEngine} 作为能力引擎门面,自动发现所有 {@link CapabilitySubEngine} 实现,
* 在所有单例 Bean 初始化完成后将子引擎注册到 CapabilityEngine。
* </p>
* <p>
* 默认注册两个子引擎:
* <ul>
* <li>{@link ToolCapabilityEngine} — 适配 ApiToolEngineFeign/OpenAPI 工具)</li>
* <li>{@link McpCapabilityEngine} — MCP 协议工具(条件注册,需 MCP SDK + McpServerConfigProvider</li>
* </ul>
* 其他子引擎(SkillCapabilityEngine / WorkflowCapabilityEngine
* 通过实现 CapabilitySubEngine 接口并声明为 Bean 即可自动接入。
* </p>
* <p>
* <b>与 ApiToolAutoConfiguration 的关系</b>
* <ul>
* <li>ApiToolAutoConfiguration 负责 ApiToolEngine 的初始化(工具发现与聚合)</li>
* <li>CapabilityAutoConfiguration 负责 CapabilityEngine 的初始化(子引擎注册与动作索引构建)</li>
* <li>两者通过 ApiToolEngine Bean 关联:ToolCapabilityEngine 依赖 ApiToolEngine</li>
* <li>初始化顺序由 Spring 容器保证:ApiToolEngine 的 SmartInitializingSingleton 先执行,
* CapabilityEngine 的 SmartInitializingSingleton 后执行</li>
* </ul>
* </p>
* <p>
* 通过 META-INF/spring/org.springframework.boot.autoconfigure.AutoConfiguration.imports 注册。
* </p>
*
* @author: mfish
* @date: 2026/07/21
*/
@Slf4j
@AutoConfiguration
public class CapabilityAutoConfiguration {
/**
* 能力引擎门面
*/
@Bean
@ConditionalOnMissingBean
public CapabilityEngine capabilityEngine() {
return new CapabilityEngine();
}
/**
* 工具能力引擎适配器(包装 ApiToolEngine 为 CapabilitySubEngine
* <p>
* 依赖 ApiToolEngine Bean,由 ApiToolAutoConfiguration 注册。
* </p>
*/
@Bean
public ToolCapabilityEngine toolCapabilityEngine(ApiToolEngine apiToolEngine) {
return new ToolCapabilityEngine(apiToolEngine);
}
/**
* MCP 能力引擎(条件注册)
* <p>
* 仅当 classpath 中存在 MCP SDK 类 且 容器中有 {@link McpServerConfigProvider} 实现时才注册。
* </p>
* <p>
* 条件说明:
* <ul>
* <li>{@code @ConditionalOnClass} — 确保 MCP SDK 在 classpath 中(mf-common-ai 已声明依赖)</li>
* <li>{@code @ConditionalOnBean(McpServerConfigProvider.class)} — 确保业务层提供了 MCP 配置数据源
* (由 mf-ai 模块的 McpServerConfigServiceImpl 实现)</li>
* </ul>
* 若条件不满足(如 mf-common-ai 独立测试或 mf-ai 模块未启动),MCP 引擎不注册,不影响其他子引擎。
* </p>
*/
@Bean
@ConditionalOnClass(name = "io.modelcontextprotocol.client.McpSyncClient")
@ConditionalOnBean(McpServerConfigProvider.class)
public McpCapabilityEngine mcpCapabilityEngine(McpServerConfigProvider configProvider) {
return new McpCapabilityEngine(configProvider);
}
/**
* 能力引擎初始化触发器:在所有单例 Bean(含所有 CapabilitySubEngine)就绪后,
* 将子引擎注册到 CapabilityEngine 并构建动作索引
* <p>
* 放在此处而非 CapabilityEngine 内,避免 CapabilityEngine 依赖 Spring 容器回调接口,保持可测试性。
* 与 ApiToolAutoConfiguration 的 apiToolEngineInitializer 模式一致。
* </p>
* <p>
* 对于 McpCapabilityEngine,在注册前先调用 {@code refresh()} 连接所有 MCP 服务器并发现工具,
* 确保 {@code getActions()} 在注册到 CapabilityEngine 时返回完整的动作列表。
* </p>
*/
@Bean
public SmartInitializingSingleton capabilityEngineInitializer(CapabilityEngine capabilityEngine,
List<CapabilitySubEngine> subEngines) {
return () -> {
log.info("[CapabilityAutoConfiguration] 发现 {} 个 CapabilitySubEngine: {}",
subEngines.size(),
subEngines.stream().map(e -> e.getEngineType().name()).toList());
// MCP 引擎需在注册前初始化(连接 MCP 服务器、发现工具)
for (CapabilitySubEngine engine : subEngines) {
if (engine instanceof McpCapabilityEngine mcpEngine) {
log.info("[CapabilityAutoConfiguration] 触发 MCP 引擎初始化");
mcpEngine.refresh();
}
}
capabilityEngine.registerSubEngines(subEngines);
};
}
}
@@ -0,0 +1,249 @@
package cn.com.mfish.common.ai.capability;
import lombok.extern.slf4j.Slf4j;
import java.util.ArrayList;
import java.util.Collections;
import java.util.LinkedHashMap;
import java.util.List;
import java.util.Map;
import java.util.concurrent.ConcurrentHashMap;
/**
* 能力引擎门面(Capability Engine
* <p>
* 组合(Composite)所有 {@link CapabilitySubEngine} 实例,向上为 Planner 提供统一的能力发现与执行入口。
* 屏蔽异构执行细节:Planner 无需关心动作来自 TOOL / SKILL / MCP / WORKFLOW 哪个引擎。
* </p>
* <p>
* <b>架构定位</b>
* <pre>
* CapabilityEngine(门面)
* │
* ┌───────────────┼───────────────┬───────────────┬───────────────┐
* │ │ │ │
* ToolCapability SkillCapability McpCapability WorkflowCapability
* (Engine) (Engine) (Engine) (Engine)
* │ │ │ │
* ApiToolEngine Java Skill MCP Client Flowable/BPMN
* (Feign+HTTP)
* </pre>
* </p>
* <p>
* <b>与现有 ApiToolEngine 的关系</b>
* <ul>
* <li>ApiToolEngine 保持不变,仍由 {@code BaseAssistant.chatWithTools()} 直接调用,
* 支撑 LLM 驱动模式(把工具交给 ChatClient</li>
* <li>CapabilityEngine 是上层抽象,ToolCapabilityEngine 适配 ApiToolEngine 暴露其动作元数据</li>
* <li>两条链路并存:LLM 驱动走 ChatClient.tools(),显式调用走 CapabilityEngine.executeAction()</li>
* </ul>
* </p>
* <p>
* <b>线程安全</b>:子引擎注册在初始化期完成,运行期只读;动作查找通过 ConcurrentHashMap 支持并发。
* </p>
*
* @author: mfish
* @date: 2026/07/21
*/
@Slf4j
public class CapabilityEngine {
/**
* 已注册的子引擎列表(按注册顺序)
*/
private final List<CapabilitySubEngine> subEngines = Collections.synchronizedList(new ArrayList<>());
/**
* 动作名 → 子引擎 的索引(供 executeAction 快速路由)
* <p>
* 初始化时构建,子引擎动态新增动作时需调用 {@link #refreshActionIndex()} 重建。
* 同名动作按子引擎注册顺序,先注册的覆盖(与 ApiToolEngine 去重策略一致)。
* </p>
*/
private volatile Map<String, CapabilitySubEngine> actionIndex = new ConcurrentHashMap<>();
/**
* 注册子引擎
* <p>
* 在 {@link CapabilityAutoConfiguration} 初始化期调用,运行期不应动态增减。
* 注册后自动重建动作索引。
* </p>
*
* @param subEngine 子引擎实例
*/
public synchronized void registerSubEngine(CapabilitySubEngine subEngine) {
if (subEngine == null) {
return;
}
subEngines.add(subEngine);
log.info("[CapabilityEngine] 注册子引擎 type={} actions={}",
subEngine.getEngineType(),
subEngine.getActions() != null ? subEngine.getActions().size() : 0);
rebuildActionIndex();
}
/**
* 批量注册子引擎
*
* @param engines 子引擎列表
*/
public synchronized void registerSubEngines(List<CapabilitySubEngine> engines) {
if (engines == null || engines.isEmpty()) {
return;
}
for (CapabilitySubEngine engine : engines) {
if (engine != null) {
subEngines.add(engine);
}
}
log.info("[CapabilityEngine] 批量注册 {} 个子引擎", subEngines.size());
rebuildActionIndex();
}
/**
* 获取所有子引擎的可用动作(并集)
* <p>
* Planner 调用此方法获取全局动作列表,可选择:
* <ul>
* <li>将动作列表注入系统提示词,供 LLM 决策</li>
* <li>将动作列表转换为 ToolCallback 传给 ChatClientLLM 驱动模式)</li>
* <li>根据动作元数据自行决策后调用 {@link #executeAction}(显式调用模式)</li>
* </ul>
* </p>
*
* @return 所有子引擎动作的并集,无子引擎时返回空列表
*/
public List<ActionDefinition> getAvailableActions() {
List<ActionDefinition> all = new ArrayList<>();
for (CapabilitySubEngine engine : subEngines) {
try {
List<ActionDefinition> actions = engine.getActions();
if (actions != null) {
all.addAll(actions);
}
} catch (Exception e) {
log.error("[CapabilityEngine] 子引擎 type={} 获取动作列表失败",
engine.getEngineType(), e);
}
}
return all;
}
/**
* 按引擎类型过滤可用动作
*
* @param engineType 引擎类型
* @return 该类型引擎的动作列表
*/
public List<ActionDefinition> getAvailableActions(EngineType engineType) {
List<ActionDefinition> filtered = new ArrayList<>();
for (CapabilitySubEngine engine : subEngines) {
if (engine.getEngineType() == engineType) {
try {
List<ActionDefinition> actions = engine.getActions();
if (actions != null) {
filtered.addAll(actions);
}
} catch (Exception e) {
log.error("[CapabilityEngine] 子引擎 type={} 获取动作列表失败", engineType, e);
}
}
}
return filtered;
}
/**
* 执行动作(统一入口)
* <p>
* 按 actionName 路由到对应子引擎执行。Planner 显式调用模式走此方法。
* </p>
*
* @param actionName 动作名称(需与 ActionDefinition.name 一致)
* @param params 动作参数
* @param ctx 执行上下文
* @return 执行结果;动作不存在时返回 failure
*/
public ActionResult executeAction(String actionName, Map<String, Object> params, ExecutionContext ctx) {
if (actionName == null || actionName.isEmpty()) {
return ActionResult.failure(EngineType.TOOL, "actionName 不能为空", 0);
}
CapabilitySubEngine engine = actionIndex.get(actionName);
if (engine == null) {
return ActionResult.failure(EngineType.TOOL,
"未找到动作: " + actionName + ",可用动作: " + actionIndex.keySet(), 0);
}
long start = System.currentTimeMillis();
try {
ActionResult result = engine.execute(actionName, params, ctx);
long elapsed = System.currentTimeMillis() - start;
// 子引擎可能未填充耗时,此处兜底
if (result.getDurationMs() <= 0) {
result.setDurationMs(elapsed);
}
log.info("[CapabilityEngine] 执行动作 name={} engine={} success={} cost={}ms",
actionName, engine.getEngineType(), result.isSuccess(), elapsed);
return result;
} catch (Exception e) {
long elapsed = System.currentTimeMillis() - start;
log.error("[CapabilityEngine] 执行动作异常 name={} engine={}",
actionName, engine.getEngineType(), e);
return ActionResult.failure(engine.getEngineType(),
"执行异常: " + e.getMessage(), elapsed);
}
}
/**
* 重建动作索引(子引擎动态新增动作后调用)
*/
public synchronized void refreshActionIndex() {
rebuildActionIndex();
}
/**
* 获取已注册的子引擎类型列表(供调试和监控)
*/
public List<EngineType> getRegisteredEngineTypes() {
List<EngineType> types = new ArrayList<>();
for (CapabilitySubEngine engine : subEngines) {
types.add(engine.getEngineType());
}
return types;
}
/**
* 获取已注册的动作总数(供调试和监控)
*/
public int getActionCount() {
return actionIndex.size();
}
/**
* 重建动作索引:遍历所有子引擎,收集动作名 → 子引擎映射
* <p>
* 同名动作按子引擎注册顺序,先注册的保留(与 ApiToolEngine 去重策略一致)。
* </p>
*/
private void rebuildActionIndex() {
Map<String, CapabilitySubEngine> newIndex = new LinkedHashMap<>();
for (CapabilitySubEngine engine : subEngines) {
try {
List<ActionDefinition> actions = engine.getActions();
if (actions == null) {
continue;
}
for (ActionDefinition action : actions) {
if (action == null || action.getName() == null) {
continue;
}
// 先注册的保留,后注册的跳过(避免同名冲突)
newIndex.putIfAbsent(action.getName(), engine);
}
} catch (Exception e) {
log.error("[CapabilityEngine] 重建索引时子引擎 type={} 获取动作失败",
engine.getEngineType(), e);
}
}
this.actionIndex = new ConcurrentHashMap<>(newIndex);
log.info("[CapabilityEngine] 动作索引重建完成,共 {} 个动作", actionIndex.size());
}
}
@@ -0,0 +1,64 @@
package cn.com.mfish.common.ai.capability;
import java.util.List;
import java.util.Map;
/**
* 能力子引擎接口
* <p>
* 屏蔽异构执行细节的统一抽象。每种能力来源(Tool / Skill / MCP / Workflow
* 实现此接口,由 {@link CapabilityEngine} 门面组合所有子引擎实例。
* </p>
* <p>
* <b>三种执行模式的关系</b>
* <ul>
* <li><b>LLM 驱动模式</b>Planner 将 {@link #getActions()} 的结果作为工具列表传给 ChatClient
* 由 LLM 决定调用哪个动作、传什么参数。TOOL 引擎目前走此模式。</li>
* <li><b>显式调用模式</b>Planner 直接通过 {@link #execute} 指定 actionName 和参数执行。
* SKILL/MCP/WORKFLOW 引擎适合此模式。</li>
* <li><b>混合模式</b>Planner 先用 getActions() 让 LLM 决策,
* 再用 execute() 执行 LLM 选择的动作。两种模式可共存。</li>
* </ul>
* </p>
* <p>
* <b>线程安全</b>:实现类需保证 getActions() 和 execute() 线程安全,
* 因为可能被多个请求线程并发调用。
* </p>
*
* @author: mfish
* @date: 2026/07/21
*/
public interface CapabilitySubEngine {
/**
* 引擎类型识别
*
* @return 引擎类型枚举(TOOL / SKILL / MCP / WORKFLOW
*/
EngineType getEngineType();
/**
* 向 Planner 暴露的可用动作元数据
* <p>
* 遵循 OpenAI Tool 格式,{@link CapabilityEngine} 将所有子引擎的返回值合并为统一列表。
* 动作名需全局唯一(建议加引擎前缀,如 {@code tool.xxx}、{@code skill.xxx})。
* </p>
*
* @return 动作定义列表,无可用动作时返回空列表
*/
List<ActionDefinition> getActions();
/**
* 执行动作
* <p>
* 显式调用入口:按 actionName 查找并执行对应动作。
* 实现类需处理 actionName 不存在的情况(返回 failure 结果,不抛异常)。
* </p>
*
* @param actionName 动作名称(需与 {@link ActionDefinition#getName()} 一致)
* @param params 动作参数(key 为参数名,value 为参数值)
* @param ctx 执行上下文(租户信息、会话ID、自定义属性)
* @return 执行结果
*/
ActionResult execute(String actionName, Map<String, Object> params, ExecutionContext ctx);
}
@@ -0,0 +1,50 @@
package cn.com.mfish.common.ai.capability;
/**
* 能力子引擎类型
* <p>
* 标识 {@link CapabilitySubEngine} 的具体来源,供 Planner 决策时区分能力来源,
* 也用于日志和调试。新增引擎类型时在此枚举扩展。
* </p>
*
* @author: mfish
* @date: 2026/07/21
*/
public enum EngineType {
/**
* 工具引擎:基于 Feign/OpenAPI 的微服务接口工具
* <p>
* 适配现有 {@code ApiToolEngine} + {@code ToolProvider} 体系,
* 将 Spring AI {@code ToolCallback} 转换为 {@link ActionDefinition}。
* </p>
*/
TOOL,
/**
* 技能引擎:Java 编写的预置技能(如文件解析、数据格式化、SQL 生成等)
* <p>
* 与 TOOL 的区别:TOOL 是远程接口的包装,SKILL 是本地 Java 逻辑的包装。
* SKILL 不依赖网络调用,执行延迟低、结果确定性强。
* </p>
*/
SKILL,
/**
* MCP 引擎:Model Context Protocol 客户端工具
* <p>
* 对接外部 MCP Server,将 MCP 工具转换为统一的 ActionDefinition。
* 当前为预留类型,待 Spring AI MCP 集成后实现。
* </p>
*/
MCP,
/**
* 工作流引擎:基于 Flowable/BPMN 的流程编排能力
* <p>
* 将工作流的"发起流程/审批/查询任务"等操作封装为 Action
* 供 Planner 在需要人工审批或复杂流程编排时调用。
* </p>
*/
WORKFLOW;
}
@@ -0,0 +1,88 @@
package cn.com.mfish.common.ai.capability;
import cn.com.mfish.common.ai.agent.TenantContext;
import lombok.Data;
import lombok.experimental.Accessors;
import java.util.HashMap;
import java.util.Map;
/**
* 能力执行上下文
* <p>
* 封装动作执行时所需的环境信息,由调用方(Planner / Executor / BaseAssistant)构建,
* 传递给 {@link CapabilitySubEngine#execute}。
* </p>
* <p>
* <b>与 Spring AI ToolContext 的关系</b>
* <ul>
* <li>ToolContext{@code Map<String, Object>})是 Spring AI 的工具执行上下文,
* 由 ChatClient 框架在调用 ToolCallback 时注入</li>
* <li>ExecutionContext 是能力引擎层的抽象,更结构化地封装租户信息和自定义属性</li>
* <li>对于 TOOL 引擎:{@link ToolCapabilityEngine} 会从 ExecutionContext 提取信息,
* 转换为 Spring AI ToolContext 传给 ToolCallback</li>
* <li>对于 SKILL/MCP/WORKFLOW 引擎:直接使用 ExecutionContext 的字段</li>
* </ul>
* </p>
*
* @author: mfish
* @date: 2026/07/21
*/
@Data
@Accessors(chain = true)
public class ExecutionContext {
/**
* 租户上下文(包含 tenantId/userId/accessToken/请求上下文快照)
* <p>
* 由请求线程捕获,供异步执行线程恢复认证上下文。
* 可为 null(如离线测试场景),各子引擎需做空值防御。
* </p>
*/
private TenantContext tenantContext;
/**
* 会话ID(用于日志追踪和会话级状态关联)
*/
private String sessionId;
/**
* 自定义属性(扩展通道,子引擎可按约定存取特定 key)
* <p>
* 常用 key(可选,非强制约定):
* <ul>
* <li>{@code requestAttributes} —— Servlet RequestAttributes</li>
* <li>{@code serverWebExchange} —— WebFlux ServerWebExchange</li>
* <li>{@code actionTraceId} —— 动作执行追踪ID</li>
* </ul>
* </p>
*/
private Map<String, Object> attributes = new HashMap<>();
/**
* 添加自定义属性
*/
public ExecutionContext addAttribute(String key, Object value) {
if (this.attributes == null) {
this.attributes = new HashMap<>();
}
this.attributes.put(key, value);
return this;
}
/**
* 获取自定义属性
*/
public Object getAttribute(String key) {
return this.attributes != null ? this.attributes.get(key) : null;
}
/**
* 从 TenantContext 快速构建执行上下文
*/
public static ExecutionContext of(TenantContext tenantContext, String sessionId) {
return new ExecutionContext()
.setTenantContext(tenantContext)
.setSessionId(sessionId);
}
}
@@ -0,0 +1,437 @@
package cn.com.mfish.common.ai.capability;
import cn.com.mfish.common.core.utils.StringUtils;
import com.alibaba.fastjson2.JSON;
import io.modelcontextprotocol.client.McpClient;
import io.modelcontextprotocol.client.McpSyncClient;
import io.modelcontextprotocol.client.transport.ServerParameters;
import io.modelcontextprotocol.client.transport.StdioClientTransport;
import io.modelcontextprotocol.json.McpJsonDefaults;
import io.modelcontextprotocol.json.McpJsonMapper;
import io.modelcontextprotocol.spec.McpClientTransport;
import io.modelcontextprotocol.spec.McpSchema;
import org.springframework.ai.mcp.client.webflux.transport.WebClientStreamableHttpTransport;
import org.springframework.ai.mcp.client.webflux.transport.WebFluxSseClientTransport;
import lombok.extern.slf4j.Slf4j;
import org.springframework.web.reactive.function.client.WebClient;
import java.time.Duration;
import java.util.ArrayList;
import java.util.Collections;
import java.util.HashMap;
import java.util.List;
import java.util.Map;
import java.util.concurrent.ConcurrentHashMap;
/**
* MCP 能力引擎
* <p>
* 实现 {@link CapabilitySubEngine} 接口,对接外部 MCP (Model Context Protocol) Server
* 将远程工具转换为统一的 {@link ActionDefinition}。
* </p>
* <p>
* <b>架构定位</b>
* <pre>
* CapabilityEngine(门面)
* │
* ├── ToolCapabilityEngine (Feign/OpenAPI)
* ├── McpCapabilityEngine (本类,MCP 协议)
* ├── [SkillCapabilityEngine] (未来扩展)
* └── [WorkflowCapabilityEngine] (未来扩展)
* </pre>
* </p>
* <p>
* <b>初始化流程</b>
* <ol>
* <li>从 {@link McpServerConfigProvider} 获取所有活跃的 MCP 服务器配置</li>
* <li>为每个配置创建 {@link McpSyncClient}
* <ul>
* <li>stdio Transport:通过 {@link StdioClientTransport} 拉起本地进程(Node.js/Python</li>
* <li>SSE Transport:通过 {@link WebFluxSseClientTransport} 连接远程 MCP 服务(已 deprecated</li>
* <li>Streamable HTTP Transport:通过 {@link WebClientStreamableHttpTransport} 连接远程 MCP 服务(MCP 2025-03-26 规范,推荐)</li>
* </ul>
* </li>
* <li>调用 {@code client.initialize()} + {@code client.listTools()} 获取工具列表</li>
* <li>将工具映射为 {@link ActionDefinition},动作名加 {@code mcp.{serverName}.} 前缀避免冲突</li>
* </ol>
* </p>
* <p>
* <b>动作名约定</b>{@code mcp.{serverName}.{toolName}}
* <br>示例:{@code mcp.filesystem.readFile}、{@code mcp.github.searchRepos}
* </p>
* <p>
* <b>线程安全</b>clients 和 actions 映射使用 ConcurrentHashMap,支持并发读取。
* 初始化在 {@link cn.com.mfish.common.ai.capability.CapabilityAutoConfiguration} 的
* SmartInitializingSingleton 中触发。
* </p>
*
* @author: mfish
* @date: 2026/07/21
*/
@Slf4j
@SuppressWarnings("deprecation")
public class McpCapabilityEngine implements CapabilitySubEngine {
/**
* stdio 传输类型标识
*/
private static final String TRANSPORT_STDIO = "stdio";
/**
* SSE 传输类型标识
*/
private static final String TRANSPORT_SSE = "sse";
/**
* Streamable HTTP 传输类型标识(MCP 2025-03-26 规范,推荐用于新服务器)
*/
private static final String TRANSPORT_STREAMABLE = "streamable";
/**
* MCP 动作名前缀
*/
private static final String MCP_ACTION_PREFIX = "mcp.";
/**
* MCP 服务器配置提供者(由业务层注入,查询数据库)
*/
private final McpServerConfigProvider configProvider;
/**
* MCP 客户端映射:serverName → McpSyncClient
*/
private final Map<String, McpSyncClient> clients = new ConcurrentHashMap<>();
/**
* 动作映射:actionName → McpActionEntry(含 client 引用和原始 toolName
*/
private final Map<String, McpActionEntry> actionRegistry = new ConcurrentHashMap<>();
/**
* 已映射的 ActionDefinition 列表(不可变快照,getActions() 直接返回)
*/
private volatile List<ActionDefinition> cachedActions = Collections.emptyList();
public McpCapabilityEngine(McpServerConfigProvider configProvider) {
this.configProvider = configProvider;
}
@Override
public EngineType getEngineType() {
return EngineType.MCP;
}
/**
* 获取所有 MCP 服务器的工具列表(并集)
* <p>
* 返回的是初始化时缓存的快照,不会实时调用 MCP 服务器。
* 如需刷新,调用 {@link #refresh()}。
* </p>
*/
@Override
public List<ActionDefinition> getActions() {
return cachedActions;
}
/**
* 执行 MCP 工具动作
* <p>
* 按 actionName 从 actionRegistry 查找对应的 McpSyncClient 和原始 toolName
* 调用 {@code client.callTool(CallToolRequest)} 执行工具。
* </p>
*/
@Override
public ActionResult execute(String actionName, Map<String, Object> params, ExecutionContext ctx) {
long start = System.currentTimeMillis();
McpActionEntry entry = actionRegistry.get(actionName);
if (entry == null) {
return ActionResult.failure(EngineType.MCP,
"未找到 MCP 动作: " + actionName, System.currentTimeMillis() - start);
}
try {
McpSchema.CallToolRequest request = new McpSchema.CallToolRequest(
entry.toolName, params != null ? params : Collections.emptyMap());
McpSchema.CallToolResult result = entry.client.callTool(request);
String output = extractTextContent(result);
return ActionResult.success(EngineType.MCP, output, System.currentTimeMillis() - start);
} catch (Exception e) {
log.error("[McpCapabilityEngine] 执行 MCP 工具失败 action={} tool={}",
actionName, entry.toolName, e);
return ActionResult.failure(EngineType.MCP,
"MCP 工具执行异常: " + e.getMessage(), System.currentTimeMillis() - start);
}
}
/**
* 初始化:连接所有 MCP 服务器,发现工具,构建动作注册表
* <p>
* 由 {@link CapabilityAutoConfiguration} 的 SmartInitializingSingleton 触发。
* 某个服务器连接失败不影响其他服务器。
* </p>
*/
public synchronized void refresh() {
if (configProvider == null) {
log.warn("[McpCapabilityEngine] 无 McpServerConfigProvider,跳过初始化");
return;
}
List<McpServerInfo> configs;
try {
configs = configProvider.getActiveServerConfigs();
} catch (Exception e) {
log.error("[McpCapabilityEngine] 获取 MCP 配置失败", e);
return;
}
if (configs == null || configs.isEmpty()) {
log.info("[McpCapabilityEngine] 无活跃的 MCP 服务器配置");
clearAll();
return;
}
// 清理旧的连接和动作
clearAll();
List<ActionDefinition> actions = new ArrayList<>();
for (McpServerInfo config : configs) {
try {
McpSyncClient client = connectServer(config);
if (client == null) {
continue;
}
clients.put(config.getServerName(), client);
// 列出工具并映射为 ActionDefinition
McpSchema.ListToolsResult toolsResult = client.listTools();
if (toolsResult == null || toolsResult.tools() == null) {
continue;
}
for (McpSchema.Tool tool : toolsResult.tools()) {
String actionName = buildActionName(config.getServerName(), tool.name());
ActionDefinition action = new ActionDefinition()
.setName(actionName)
.setDescription(tool.description())
.setInputSchema(serializeSchema(tool.inputSchema()))
.setEngineType(EngineType.MCP)
.setServiceId(config.getServerName());
actions.add(action);
actionRegistry.put(actionName, new McpActionEntry(client, tool.name()));
log.info("[McpCapabilityEngine] 注册 MCP 工具 server={} tool={} action={}",
config.getServerName(), tool.name(), actionName);
}
log.info("[McpCapabilityEngine] MCP 服务器 {} 连接成功,注册 {} 个工具",
config.getServerName(), toolsResult.tools().size());
} catch (Exception e) {
log.error("[McpCapabilityEngine] MCP 服务器 {} 连接失败", config.getServerName(), e);
}
}
this.cachedActions = Collections.unmodifiableList(actions);
log.info("[McpCapabilityEngine] 初始化完成,共 {} 个 MCP 服务器,{} 个工具",
clients.size(), actions.size());
}
/**
* 连接单个 MCP 服务器(支持 stdio / SSE / Streamable HTTP 三种 Transport
*/
private McpSyncClient connectServer(McpServerInfo config) {
String transportType = config.getTransportType();
if (TRANSPORT_STDIO.equalsIgnoreCase(transportType)) {
return connectStdio(config);
} else if (TRANSPORT_SSE.equalsIgnoreCase(transportType)) {
return connectSse(config);
} else if (TRANSPORT_STREAMABLE.equalsIgnoreCase(transportType)) {
return connectStreamable(config);
} else {
log.warn("[McpCapabilityEngine] 不支持的传输类型: {} server={}", transportType, config.getServerName());
return null;
}
}
/**
* stdio Transport:拉起本地进程
*/
private McpSyncClient connectStdio(McpServerInfo config) {
if (StringUtils.isEmpty(config.getCommand())) {
log.warn("[McpCapabilityEngine] stdio 模式缺少 command 参数 server={}", config.getServerName());
return null;
}
ServerParameters.Builder builder = ServerParameters.builder(config.getCommand());
// 解析 argsJSON 数组字符串 → List<String>
if (StringUtils.isNotEmpty(config.getArgs())) {
try {
List<String> argsList = JSON.parseArray(config.getArgs(), String.class);
builder.args(argsList);
} catch (Exception e) {
log.warn("[McpCapabilityEngine] 解析 args 失败 server={} args={}",
config.getServerName(), config.getArgs(), e);
}
}
// 解析 envJSON 对象字符串 → Map<String, String>
if (StringUtils.isNotEmpty(config.getEnv())) {
try {
Map<String, String> envMap = JSON.parseObject(config.getEnv(), Map.class);
builder.env(envMap);
} catch (Exception e) {
log.warn("[McpCapabilityEngine] 解析 env 失败 server={} env={}",
config.getServerName(), config.getEnv(), e);
}
}
// MCP SDK 2.0.0StdioClientTransport 构造函数需要 McpJsonMapper 参数
McpJsonMapper jsonMapper = McpJsonDefaults.getMapper();
StdioClientTransport transport = new StdioClientTransport(builder.build(), jsonMapper);
return createAndInitialize(transport, config.getServerName());
}
/**
* SSE Transport:连接远程 MCP 服务
*/
private McpSyncClient connectSse(McpServerInfo config) {
if (StringUtils.isEmpty(config.getSseUrl())) {
log.warn("[McpCapabilityEngine] SSE 模式缺少 sseUrl 参数 server={}", config.getServerName());
return null;
}
WebClient.Builder webClientBuilder = WebClient.builder()
.baseUrl(config.getSseUrl());
// 注入认证 Token
if (StringUtils.isNotEmpty(config.getAuthToken())) {
webClientBuilder.defaultHeader("Authorization", "Bearer " + config.getAuthToken());
}
WebFluxSseClientTransport.Builder transportBuilder = WebFluxSseClientTransport.builder(webClientBuilder);
if (StringUtils.isNotEmpty(config.getSseEndpoint())) {
transportBuilder.sseEndpoint(config.getSseEndpoint());
}
WebFluxSseClientTransport transport = transportBuilder.build();
return createAndInitialize(transport, config.getServerName());
}
/**
* Streamable HTTP Transport:连接远程 MCP 服务(MCP 2025-03-26 规范,推荐用于新服务器)
* <p>
* 与 SSE 的区别:
* <ul>
* <li>单一 /mcp 端点,请求可流式响应也可普通响应</li>
* <li>支持会话恢复(resumableStreams</li>
* <li>支持协议版本协商(supportedProtocolVersions</li>
* <li>是 MCP 2025-03-26 规范的标准 transportSSE 已标记 deprecated</li>
* </ul>
* </p>
* <p>
* 字段复用:sseUrl 作为基础 URLsseEndpoint 作为端点路径(通常为 /mcp)。
* </p>
*/
private McpSyncClient connectStreamable(McpServerInfo config) {
if (StringUtils.isEmpty(config.getSseUrl())) {
log.warn("[McpCapabilityEngine] streamable 模式缺少 sseUrl 参数 server={}", config.getServerName());
return null;
}
WebClient.Builder webClientBuilder = WebClient.builder()
.baseUrl(config.getSseUrl());
// 注入认证 Token
if (StringUtils.isNotEmpty(config.getAuthToken())) {
webClientBuilder.defaultHeader("Authorization", "Bearer " + config.getAuthToken());
}
WebClientStreamableHttpTransport.Builder transportBuilder = WebClientStreamableHttpTransport.builder(webClientBuilder);
if (StringUtils.isNotEmpty(config.getSseEndpoint())) {
transportBuilder.endpoint(config.getSseEndpoint());
}
WebClientStreamableHttpTransport transport = transportBuilder.build();
return createAndInitialize(transport, config.getServerName());
}
/**
* 创建 McpSyncClient 并初始化连接
*/
private McpSyncClient createAndInitialize(McpClientTransport transport, String serverName) {
McpSyncClient client = McpClient.sync(transport)
.requestTimeout(Duration.ofSeconds(60))
.build();
client.initialize();
log.info("[McpCapabilityEngine] MCP 客户端已初始化 server={}", serverName);
return client;
}
/**
* 构建 MCP 动作名:mcp.{serverName}.{toolName}
*/
private String buildActionName(String serverName, String toolName) {
return MCP_ACTION_PREFIX + serverName + "." + toolName;
}
/**
* 从 CallToolResult 提取文本内容
*/
private String extractTextContent(McpSchema.CallToolResult result) {
if (result == null || result.content() == null || result.content().isEmpty()) {
return "";
}
StringBuilder sb = new StringBuilder();
for (McpSchema.Content content : result.content()) {
if (content instanceof McpSchema.TextContent tc) {
sb.append(tc.text());
} else {
// 非 TextContent 类型,序列化为 JSON
sb.append(JSON.toJSONString(content));
}
sb.append("\n");
}
return sb.toString().trim();
}
/**
* 将 MCP Tool 的 inputSchema 序列化为 JSON 字符串
*/
private String serializeSchema(Object inputSchema) {
if (inputSchema == null) {
return "{}";
}
try {
return JSON.toJSONString(inputSchema);
} catch (Exception e) {
log.warn("[McpCapabilityEngine] 序列化 inputSchema 失败", e);
return "{}";
}
}
/**
* 清理所有连接和动作注册表
*/
private void clearAll() {
for (Map.Entry<String, McpSyncClient> entry : clients.entrySet()) {
try {
entry.getValue().close();
} catch (Exception e) {
log.warn("[McpCapabilityEngine] 关闭 MCP 客户端失败 server={}", entry.getKey(), e);
}
}
clients.clear();
actionRegistry.clear();
cachedActions = Collections.emptyList();
}
/**
* MCP 动作注册表条目:持有 McpSyncClient 引用和原始 toolName
*/
private static class McpActionEntry {
final McpSyncClient client;
final String toolName;
McpActionEntry(McpSyncClient client, String toolName) {
this.client = client;
this.toolName = toolName;
}
}
}
@@ -0,0 +1,27 @@
package cn.com.mfish.common.ai.capability;
import java.util.List;
/**
* MCP 服务器配置提供者接口
* <p>
* 跨模块解耦:mf-common-ai 模块的 {@link McpCapabilityEngine} 通过此接口获取数据库中的 MCP 配置,
* 由 mf-ai 业务模块实现(查询 ai_mcp_server_config 表)。
* </p>
* <p>
* 若 Spring 容器中无此接口的实现 Bean(如 mf-common-ai 独立测试场景),
* {@link McpCapabilityEngine} 不报错,返回空动作列表。
* </p>
*
* @author: mfish
* @date: 2026/07/21
*/
public interface McpServerConfigProvider {
/**
* 获取所有状态为"正常"的 MCP 服务器配置
*
* @return MCP 服务器配置列表,无配置时返回空列表
*/
List<McpServerInfo> getActiveServerConfigs();
}
@@ -0,0 +1,64 @@
package cn.com.mfish.common.ai.capability;
import lombok.Data;
import lombok.experimental.Accessors;
/**
* MCP 服务器配置信息(mf-common-ai 模块内的 DTO
* <p>
* 与 {@code McpServerConfig} 实体类解耦:mf-common-ai 不依赖 mf-ai-api 模块,
* 通过 {@link McpServerConfigProvider} 接口由业务层注入配置数据。
* </p>
*
* @author: mfish
* @date: 2026/07/21
*/
@Data
@Accessors(chain = true)
public class McpServerInfo {
/**
* MCP 服务器名称(唯一标识)
*/
private String serverName;
/**
* 传输类型:stdio / sse / streamable
* <ul>
* <li>stdio — 拉起本地进程(Node.js/Python</li>
* <li>sse — SSE TransportMCP 2024-11-05 规范,已 deprecated</li>
* <li>streamable — Streamable HTTP TransportMCP 2025-03-26 规范,推荐)</li>
* </ul>
*/
private String transportType;
/**
* stdio 模式启动命令(如 node / python
*/
private String command;
/**
* stdio 模式参数(JSON 数组字符串,如 ["server.js","--port","3000"]
*/
private String args;
/**
* stdio 模式环境变量(JSON 对象字符串)
*/
private String env;
/**
* 远程 MCP 服务基础 URLsse / streamable 模式通用,如 https://mcp.example.com
*/
private String sseUrl;
/**
* 远程 MCP 服务端点路径(sse / streamable 模式通用,如 /sse 或 /mcp
*/
private String sseEndpoint;
/**
* 认证 TokenBearer
*/
private String authToken;
}
@@ -0,0 +1,175 @@
package cn.com.mfish.common.ai.capability;
import cn.com.mfish.common.ai.agent.TenantContext;
import cn.com.mfish.common.ai.engine.ApiToolEngine;
import cn.com.mfish.common.core.constants.RPCConstants;
import cn.com.mfish.common.core.utils.AuthInfoUtils;
import cn.com.mfish.common.core.utils.StringUtils;
import com.alibaba.fastjson2.JSON;
import lombok.extern.slf4j.Slf4j;
import org.springframework.ai.chat.model.ToolContext;
import org.springframework.ai.tool.ToolCallback;
import org.springframework.ai.tool.definition.ToolDefinition;
import org.springframework.web.context.request.RequestAttributes;
import org.springframework.web.server.ServerWebExchange;
import java.util.ArrayList;
import java.util.HashMap;
import java.util.List;
import java.util.Map;
/**
* 工具能力引擎适配器
* <p>
* 将现有 {@link ApiToolEngine} 包装为 {@link CapabilitySubEngine}
* 使 Feign/OpenAPI 工具能纳入 CapabilityEngine 的统一能力发现与执行体系。
* </p>
* <p>
* <b>适配逻辑</b>
* <ul>
* <li>{@link #getActions()}:遍历 ApiToolEngine 所有 ToolCallback
* 将 {@code ToolDefinition}name/description/inputSchema)转换为 {@link ActionDefinition}</li>
* <li>{@link #execute}:按 actionName 查找 ToolCallback
* 将 params Map 序列化为 JSON 字符串,构建 Spring AI ToolContext
* 调用 {@code ToolCallback.call(toolInput, toolContext)}</li>
* </ul>
* </p>
* <p>
* <b>与 LLM 驱动模式的关系</b>
* <ul>
* <li>LLM 驱动模式:BaseAssistant 直接调用 {@code apiToolEngine.getToolCallbackProvider(serviceIds)}
* 传给 ChatClient,由 LLM 决定调用哪个工具。此模式不经过 ToolCapabilityEngine。</li>
* <li>显式调用模式:Planner 通过 {@code capabilityEngine.executeAction(actionName, params, ctx)}
* 直接执行指定动作。此模式经过 ToolCapabilityEngine。</li>
* <li>两种模式共享同一套 ToolCallback 实例,执行结果一致。</li>
* </ul>
* </p>
*
* @author: mfish
* @date: 2026/07/21
*/
@Slf4j
public class ToolCapabilityEngine implements CapabilitySubEngine {
private final ApiToolEngine apiToolEngine;
public ToolCapabilityEngine(ApiToolEngine apiToolEngine) {
this.apiToolEngine = apiToolEngine;
}
@Override
public EngineType getEngineType() {
return EngineType.TOOL;
}
/**
* 将 ApiToolEngine 中所有 ToolCallback 转换为 ActionDefinition
* <p>
* 动作名直接使用 ToolCallback 的原始 name(与 LLM 驱动模式保持一致),
* 跨服务同名动作由 ApiToolEngine 去重策略保证唯一。
* serviceId 字段标识动作来源服务,便于 Planner 按服务过滤。
* </p>
*/
@Override
public List<ActionDefinition> getActions() {
List<ActionDefinition> actions = new ArrayList<>();
for (String serviceId : apiToolEngine.getAllServiceIds()) {
org.springframework.ai.tool.ToolCallbackProvider provider =
apiToolEngine.getToolCallbackProvider(serviceId);
if (provider == null) {
continue;
}
org.springframework.ai.tool.ToolCallback[] callbacks = provider.getToolCallbacks();
if (callbacks == null) {
continue;
}
for (org.springframework.ai.tool.ToolCallback tc : callbacks) {
try {
ToolDefinition td = tc.getToolDefinition();
actions.add(new ActionDefinition()
.setName(td.name())
.setDescription(td.description())
.setInputSchema(td.inputSchema())
.setEngineType(EngineType.TOOL)
.setServiceId(serviceId));
} catch (Exception e) {
log.warn("[ToolCapabilityEngine] 转换动作失败 serviceId={}", serviceId, e);
}
}
}
return actions;
}
/**
* 执行工具动作
* <p>
* 通过 ApiToolEngine 查找 actionName 对应的 ToolCallback
* 将 params 序列化为 JSON 字符串作为 toolInput
* 从 ExecutionContext 提取租户信息构建 ToolContext
* 调用 {@code ToolCallback.call(toolInput, toolContext)}。
* </p>
*/
@Override
public ActionResult execute(String actionName, Map<String, Object> params, ExecutionContext ctx) {
long start = System.currentTimeMillis();
ToolCallback callback = apiToolEngine.findToolCallback(actionName);
if (callback == null) {
return ActionResult.failure(EngineType.TOOL,
"未找到工具动作: " + actionName, System.currentTimeMillis() - start);
}
// 构建 Spring AI ToolContext(复用 BaseAssistant 的 context key 约定)
ToolContext toolContext = buildToolContext(ctx);
// params 序列化为 JSON 字符串
String toolInput = (params == null || params.isEmpty())
? "{}"
: JSON.toJSONString(params);
try {
String result = callback.call(toolInput, toolContext);
return ActionResult.success(EngineType.TOOL, result, System.currentTimeMillis() - start);
} catch (Exception e) {
log.error("[ToolCapabilityEngine] 执行工具动作失败 name={} input={}",
actionName, toolInput, e);
return ActionResult.failure(EngineType.TOOL,
"工具执行异常: " + e.getMessage(), System.currentTimeMillis() - start);
}
}
/**
* 从 ExecutionContext 构建 Spring AI ToolContext
* <p>
* 复用 {@code BaseAssistant.buildToolContextFromSnapshot} 的 key 约定
* RPCConstants.REQ_* 常量),保证 FeignToolCallback/HttpToolCallback
* 的内部参数自动填充逻辑一致。
* </p>
*/
private ToolContext buildToolContext(ExecutionContext ctx) {
Map<String, Object> contextMap = new HashMap<>();
if (ctx != null && ctx.getTenantContext() != null) {
TenantContext tc = ctx.getTenantContext();
contextMap.put(RPCConstants.REQ_USER_ID, tc.getUserId() != null ? tc.getUserId() : "");
contextMap.put(RPCConstants.REQ_TENANT_ID,
tc.getTenantId() != null ? tc.getTenantId() : AuthInfoUtils.SUPER_TENANT_ID);
contextMap.put(RPCConstants.REQ_ORIGIN, RPCConstants.AI);
if (StringUtils.isNotEmpty(tc.getAccessToken())) {
contextMap.put(RPCConstants.REQ_TOKEN, tc.getAccessToken());
}
if (tc.getRequestAttributes() instanceof RequestAttributes ra) {
contextMap.put(RPCConstants.REQ_REQUEST_ATTRIBUTES, ra);
}
if (tc.getServerWebExchange() instanceof ServerWebExchange swe) {
contextMap.put(RPCConstants.REQ_SERVER_WEB_EXCHANGE, swe);
}
} else {
// 无 TenantContext 时填充默认值,避免内部参数解析失败
contextMap.put(RPCConstants.REQ_USER_ID, "");
contextMap.put(RPCConstants.REQ_TENANT_ID, AuthInfoUtils.SUPER_TENANT_ID);
contextMap.put(RPCConstants.REQ_ORIGIN, RPCConstants.AI);
}
return new ToolContext(contextMap);
}
}
@@ -186,4 +186,53 @@ public class ApiToolEngine {
public ToolCallbackProvider getToolCallbackProvider(Set<String> serviceIds) {
return ToolCallbackProvider.from(getToolCallbacks(serviceIds));
}
/**
* 获取所有已注册的微服务ID
* <p>
* 供 {@code ToolCapabilityEngine} 枚举所有服务的工具,将其暴露为 ActionDefinition。
* </p>
*
* @return 已注册的 serviceId 集合(不可变快照)
*/
public Set<String> getAllServiceIds() {
return Set.copyOf(toolsByService.keySet());
}
/**
* 获取所有服务的全部工具(跨服务,同名工具按 serviceId 顺序去重)
* <p>
* 供 {@code ToolCapabilityEngine} 构建全局动作列表。
* 去重策略与 {@link #getToolCallbacks(Collection)} 一致:先注册的保留。
* </p>
*
* @return 全部工具列表(去重后的快照)
*/
public List<ToolCallback> getAllToolCallbacks() {
return getToolCallbacks(toolsByService.keySet());
}
/**
* 按动作名查找对应的 ToolCallback
* <p>
* 供 {@code ToolCapabilityEngine.execute()} 通过 actionName 路由到具体工具。
* 遍历所有服务,返回首个名称匹配的 ToolCallback。
* </p>
*
* @param actionName 动作名(即 ToolCallback.getToolDefinition().name()
* @return 匹配的 ToolCallback,未找到时返回 null
*/
public ToolCallback findToolCallback(String actionName) {
if (actionName == null || actionName.isEmpty()) {
return null;
}
for (List<ToolCallback> callbacks : toolsByService.values()) {
for (ToolCallback tc : callbacks) {
if (actionName.equals(tc.getToolDefinition().name())) {
return tc;
}
}
}
return null;
}
}
@@ -0,0 +1,268 @@
package cn.com.mfish.common.ai.memory;
import cn.com.mfish.common.ai.agent.TenantContext;
import java.util.List;
import java.util.Map;
/**
* 会话级记忆容器(Conversation Memory
* <p>
* 三驾马车驱动的 AI 架构 —— Memory 模块的核心抽象。
* 每个 sessionId 对应一个 ConversationMemory 实例,聚合会话内 Spring AI 管不到的两类上下文:
* </p>
* <ol>
* <li><b>Vars Context(变量上下文)</b>:低代码平台运行实例中的局部变量、用户凭证(Token)、
* 租户上下文等。统一收纳原本散落在 TenantContext / ToolContext / 业务变量中的数据,
* 通过 {@link #updateVariable} / {@link #getVariable} 读写,{@link #getVarsContext}
* 供 Planner 一次性快照。</li>
* <li><b>Document Context(文档上下文)</b>:由 FileParseService 解析后的临时文件块(Chunk)。
* 通过 {@link #addDocumentChunks} 批量注入,{@link #getDocumentChunks} 读取,
* {@link #getDocumentContext} 直接产出 LLM 可读的拼接文本。</li>
* </ol>
* <p>
* <b>为什么不管理 Short-term MemoryChat 历史)?</b>
* Spring AI 的 {@code ChatMemory} + {@code MessageChatMemoryAdvisor} 已经完整覆盖
* "对话历史存储 + 自动注入 LLM 上下文" 两个职责,本模块不再重复实现,避免双写不一致。
* Planner / Assistant 若需要读取对话历史,直接注入 Spring AI 的 {@code ChatMemory} 即可。
* </p>
* <p>
* 接口设计原则:
* <ul>
* <li><b>会话隔离</b>:所有方法以当前实例绑定的 sessionId 为作用域,不跨会话泄漏</li>
* <li><b>两段独立</b>Vars 和 Documents 互不耦合,可分别清理({@link #clearVars} / {@link #clearDocuments}</li>
* <li><b>上下文拼接</b>{@link #getSystemContext()} 一次性产出 LLM 系统提示词所需的全部上下文片段</li>
* <li><b>不耦合存储</b>:本接口只定义读写契约,底层存储由实现类决定(内存/Redis/DB</li>
* </ul>
* </p>
*
* @author: mfish
* @date: 2026/07/20
*/
public interface ConversationMemory {
/**
* 获取当前会话 ID
*/
String getSessionId();
// ==================== Vars Context ====================
/**
* 更新一个变量到 Vars Context
* <p>
* 若变量已存在则覆盖。常见用法:
* </p>
* <ul>
* <li>低代码平台运行时写入局部变量(如 `userId` / `formId` / `businessKey`</li>
* <li>认证拦截器写入用户凭证(`accessToken` / `refreshToken`</li>
* <li>编排流程写入中间结果(如 Planner 产出的 `currentPlan` / Executor 写入的 `stepResults`</li>
* </ul>
*
* @param key 变量名(建议使用命名空间前缀避免冲突,如 `sys.userId` / `biz.formId`
* @param value 变量值(可为 null,表示移除该变量)
*/
void updateVariable(String key, Object value);
/**
* 读取一个变量
*
* @param key 变量名
* @return 变量值;不存在返回 null
*/
Object getVariable(String key);
/**
* 读取一个变量并按指定类型转换
*
* @param key 变量名
* @param targetType 目标类型
* @param <T> 泛型
* @return 转换后的值;不存在或类型不匹配返回 null
*/
<T> T getVariable(String key, Class<T> targetType);
/**
* 获取 Vars Context 的快照(只读视图)
* <p>
* Planner 拼接上下文时调用,避免逐个 getVariable 的开销。
* 返回的 Map 不可变,修改需通过 {@link #updateVariable}。
* </p>
*
* @return 变量快照(不可变 Map
*/
Map<String, Object> getVarsContext();
/**
* 绑定租户上下文到 Vars Context
* <p>
* 将 TenantContext 的 5 个字段(tenantId / userId / accessToken / requestAttributes /
* serverWebExchange)按约定的 key 写入 Vars Context,供后续工具调用恢复认证态。
* </p>
* <p>
* 这是对原 AgentRuntime.captureTenantContext + BaseAssistant.buildToolContext 流程的统一收纳:
* 原本散落在两处的"捕获-传递-解包"逻辑,现在统一为 Memory 的一次 bind 调用。
* </p>
*
* @param tenantContext 租户上下文快照(null 表示清除绑定的租户信息)
*/
void bindTenantContext(TenantContext tenantContext);
/**
* 清空 Vars Context
*/
void clearVars();
// ==================== Document Context ====================
/**
* 批量添加文档块到 Document Context
* <p>
* FileParseService 解析文件后调用此方法注入。重复添加相同 chunkId 的块会被去重。
* </p>
*
* @param chunks 文档块列表(不能为 null,可为空列表)
*/
void addDocumentChunks(List<DocumentChunk> chunks);
/**
* 获取当前会话的全部文档块
* <p>
* 返回的是副本或不可变视图,外部修改不影响内部状态。
* </p>
*
* @return 文档块列表(可能为空,不会为 null)
*/
List<DocumentChunk> getDocumentChunks();
/**
* 获取 Document Context 拼接后的文本
* <p>
* 按 chunkIndex 顺序拼接所有文档块内容,每个块带文件名分隔符。
* Planner / Assistant 注入 LLM 上下文时调用此方法,无需自己拼接。
* </p>
* <p>
* <b>全量注入语义</b>:本方法返回当前会话的全部文档块内容,适合文档量较小的场景。
* 文档量较大时(几十份 PDF / 大型代码库),建议改用 {@link #searchDocumentContext} 按相关性检索,
* 避免 LLM 上下文窗口爆炸。
* </p>
* <p>
* 示例输出:
* </p>
* <pre>
* === 文件: example.docx ===
* 文件内容...
* === 文件结束: example.docx ===
* </pre>
*
* @return 拼接后的文本;无文档时返回空字符串
*/
String getDocumentContext();
/**
* 按相关性检索文档块(向量库扩展点)
* <p>
* <b>默认实现(内存版)</b>:忽略 query,返回全部文档块(截断到 topK)。
* 这等价于"全量注入",保持与 {@link #getDocumentChunks()} 一致的行为。
* </p>
* <p>
* <b>向量库版实现</b>:将 query 文本向量化,按余弦相似度返回最相关的 topK 个 chunk。
* 这是 RAG(检索增强生成)的核心入口,让 LLM 只看到最相关的文档片段而非全部文档。
* </p>
* <p>
* 调用方建议:
* <ul>
* <li>文档量小(< 5 个文件或 < 20K 字符):用 {@link #getDocumentChunks()} 全量注入</li>
* <li>文档量大:用 {@code searchDocuments(prompt, 5)} 检索 topK 个最相关 chunk</li>
* </ul>
* </p>
*
* @param query 检索查询文本(通常是用户 prompt 或当前任务描述);内存版可忽略此参数
* @param topK 返回的最大 chunk 数;&lt;= 0 表示不限制(返回全部)
* @return 检索到的文档块列表(按相关性或 chunkIndex 排序);可能为空,不会为 null
*/
default List<DocumentChunk> searchDocuments(String query, int topK) {
List<DocumentChunk> all = getDocumentChunks();
if (topK <= 0 || all.size() <= topK) {
return all;
}
return new java.util.ArrayList<>(all.subList(0, topK));
}
/**
* 按相关性检索并拼接为 LLM 上下文文本(向量库扩展点)
* <p>
* 默认实现:调用 {@link #searchDocuments(query, topK)} 检索,然后按
* {@link #getDocumentContext()} 的格式拼接。
* </p>
* <p>
* 向量库版可重写此方法以使用更高效的批量检索 API,或在拼接时附加相关性分数。
* </p>
*
* @param query 检索查询文本;内存版可忽略
* @param topK 返回的最大 chunk 数;&lt;= 0 表示不限制
* @return 拼接后的文本;无文档或检索无结果时返回空字符串
*/
default String searchDocumentContext(String query, int topK) {
List<DocumentChunk> chunks = searchDocuments(query, topK);
if (chunks.isEmpty()) {
return "";
}
StringBuilder sb = new StringBuilder();
for (DocumentChunk chunk : chunks) {
String fileName = chunk.getFileName() != null ? chunk.getFileName() : chunk.getFileKey();
sb.append("=== 文件: ").append(fileName).append(" ===\n");
sb.append(chunk.getContent() != null ? chunk.getContent() : "");
sb.append("\n=== 文件结束: ").append(fileName).append(" ===\n\n");
}
return sb.toString();
}
/**
* 清空 Document Context
*/
void clearDocuments();
// ==================== 上下文拼接 ====================
/**
* 一次性产出 LLM 系统提示词所需的全部上下文
* <p>
* Planner 在构建 prompt 时调用此方法,获取 Vars Context 和 Document Context 的拼接文本,
* 注入到 system prompt 中。Short-term MemoryChat 历史)由 Spring AI
* MessageChatMemoryAdvisor 自动注入,不在此方法产出范围内。
* </p>
* <p>
* 输出格式示例:
* </p>
* <pre>
* 【用户上下文】
* 租户: tenant_001
* 用户: user_123
*
* 【业务变量】
* formId: 1001
* businessKey: ORDER-2026-001
*
* 【上传文件】
* === 文件: example.docx ===
* 文件内容...
* === 文件结束: example.docx ===
* </pre>
*
* @return 系统上下文文本;无任何上下文时返回空字符串
*/
String getSystemContext();
/**
* 清空所有记忆(Vars + Documents
* <p>
* 用于会话彻底销毁场景。会话本身的存储条目由 {@link ConversationMemoryStore} 管理。
* </p>
* <p>
* 注意:此方法不清空 Spring AI 的 ChatMemory(对话历史),那部分由 Spring AI
* 自己管理,会话结束时业务方需另行调用 {@code ChatMemory.clear(sessionId)}。
* </p>
*/
void clearAll();
}
@@ -0,0 +1,85 @@
package cn.com.mfish.common.ai.memory;
import java.util.List;
/**
* 会话记忆存储(Conversation Memory Store
* <p>
* Memory 模块的入口,负责按 sessionId 创建、获取、销毁 {@link ConversationMemory} 实例。
* 对外暴露类似 Map 的语义,但底层可以是内存 / Redis / 数据库。
* </p>
* <p>
* 使用方式:
* </p>
* <pre>
* // 在请求入口获取(或创建)会话 Memory
* ConversationMemory memory = memoryStore.getOrCreate(sessionId);
*
* // 注入上下文
* memory.bindTenantContext(tenantContext);
* memory.addDocumentChunks(chunks);
* memory.updateVariable("biz.formId", 1001);
*
* // Planner / Assistant 读取上下文
* String systemContext = memory.getSystemContext();
*
* // 会话结束后清理
* memoryStore.remove(sessionId);
* </pre>
* <p>
* 设计为 Spring Bean,由 {@code MemoryAutoConfiguration} 注册。
* Controller / Assistant / Runtime 通过依赖注入获取。
* </p>
*
* @author: mfish
* @date: 2026/07/20
*/
public interface ConversationMemoryStore {
/**
* 获取或创建会话 Memory
* <p>
* 首次访问某 sessionId 时创建新实例并缓存;后续访问返回同一实例。
* 实例的生命周期由实现类管理(如基于 TTL 自动过期)。
* </p>
*
* @param sessionId 会话ID(不能为 null 或空)
* @return 会话 Memory 实例
*/
ConversationMemory getOrCreate(String sessionId);
/**
* 获取会话 Memory(不创建)
*
* @param sessionId 会话ID
* @return 会话 Memory 实例;不存在返回 null
*/
ConversationMemory get(String sessionId);
/**
* 移除会话 Memory 并清理其内部状态
* <p>
* 触发 {@link ConversationMemory#clearAll()},并从存储中移除。
* 用于会话主动关闭、用户登出等场景。
* </p>
*
* @param sessionId 会话ID
* @return 被移除的实例;不存在返回 null
*/
ConversationMemory remove(String sessionId);
/**
* 获取当前活跃的会话 ID 列表
* <p>
* 主要用于监控、调试、定期清理任务。
* </p>
*
* @return 会话 ID 列表(不可变)
*/
List<String> listSessionIds();
/**
* 当前存储的会话数量
*/
int size();
}
@@ -0,0 +1,79 @@
package cn.com.mfish.common.ai.memory;
import lombok.Data;
import lombok.experimental.Accessors;
/**
* 文档块(Document Chunk
* <p>
* 由 FileParseService 解析文件后产出的最小内容单元。一个文件可能产生多个 Chunk,
* 也可能只有一个 Chunk(小文件)。Memory 模块按 Chunk 维度管理文档上下文,
* 便于后续扩展:按相关性检索(RAG)、按 token 上限分片注入等。
* </p>
* <p>
* 设计为不可变值对象:构造后内容不再修改,需要更新时替换整个 Chunk。
* </p>
*
* @author: mfish
* @date: 2026/07/20
*/
@Data
@Accessors(chain = true)
public class DocumentChunk {
/**
* 块ID(建议使用 UUID 或 fileKey + 序号)
*/
private String chunkId;
/**
* 来源文件的 fileKeyStorage 服务的唯一标识)
*/
private String fileKey;
/**
* 来源文件名(用于 LLM 提示词展示)
*/
private String fileName;
/**
* 文件 MIME 类型(用于 LLM 提示词展示和后续分流处理)
*/
private String fileType;
/**
* 块文本内容(已提取为纯文本,可直接注入 LLM 上下文)
*/
private String content;
/**
* 在原文件中的序号(从 0 开始;单文件单块时为 0)
*/
private int chunkIndex;
/**
* 块字符数(便于 token 估算和上下文裁剪)
*/
private int length;
/**
* 文本向量(embedding),用于向量库相似度检索
* <p>
* 维度由 embedding model 决定(如 OpenAI text-embedding-3-small 为 1536 维)。
* </p>
* <p>
* <b>使用约定</b>
* <ul>
* <li>内存版 Memory:不使用此字段(始终为 null),searchDocuments 走全量返回</li>
* <li>向量库版 Memory:由 Store 在写入时调用 embedding model 计算并存储,
* chunk 对象上的此字段可为 null(向量库内部维护真正的向量索引)</li>
* <li>调试场景:可手动填充此字段用于离线分析</li>
* </ul>
* </p>
* <p>
* 类型为 {@code float[]} 而非 {@code List<Float>},避免装箱开销;
* Lombok {@code @Data} 会基于内容生成 equals/hashCode,正常使用无影响。
* </p>
*/
private float[] embedding;
/**
* 解析时间戳(毫秒)
*/
private long parsedAt;
public DocumentChunk() {
this.parsedAt = System.currentTimeMillis();
}
}
@@ -0,0 +1,279 @@
package cn.com.mfish.common.ai.memory;
import cn.com.mfish.common.ai.agent.TenantContext;
import lombok.extern.slf4j.Slf4j;
import java.util.ArrayList;
import java.util.Collections;
import java.util.LinkedHashMap;
import java.util.List;
import java.util.Map;
import java.util.Objects;
import java.util.concurrent.ConcurrentHashMap;
/**
* 基于 JVM 内存的 {@link ConversationMemory} 默认实现
* <p>
* 两段记忆的存储策略:
* <ul>
* <li><b>Vars Context</b>:使用 {@link ConcurrentHashMap} 存储变量,支持并发读写。
* 值为 Object,业务方自行约定类型契约。</li>
* <li><b>Document Context</b>:使用 {@link List} + 同步锁保护,因为文档块通常批量写入、
* 顺序读出。按 chunkId 去重。</li>
* </ul>
* </p>
* <p>
* <b>不管理 Short-term Memory</b>:对话历史由 Spring AI 的 {@code ChatMemory} +
* {@code MessageChatMemoryAdvisor} 自动管理,本类不涉及,避免双写不一致。
* </p>
* <p>
* 线程安全:本类所有方法线程安全,可被多个线程(请求线程 + 异步编排线程)并发访问。
* 实例由 {@link InMemoryConversationMemoryStore} 按 sessionId 单例化,无需关心实例级并发。
* </p>
*
* @author: mfish
* @date: 2026/07/20
*/
@Slf4j
public class InMemoryConversationMemory implements ConversationMemory {
/**
* Vars Context 中租户上下文的标准 key 前缀
*/
private static final String TENANT_KEY_PREFIX = "sys.tenant.";
/** tenantId 在 Vars Context 中的 key */
public static final String KEY_TENANT_ID = "sys.tenant.tenantId";
/** userId 在 Vars Context 中的 key */
public static final String KEY_USER_ID = "sys.tenant.userId";
/** accessToken 在 Vars Context 中的 key */
public static final String KEY_ACCESS_TOKEN = "sys.tenant.accessToken";
/** requestAttributes 在 Vars Context 中的 keyServlet */
public static final String KEY_REQUEST_ATTRIBUTES = "sys.tenant.requestAttributes";
/** serverWebExchange 在 Vars Context 中的 keyWebFlux */
public static final String KEY_SERVER_WEB_EXCHANGE = "sys.tenant.serverWebExchange";
private final String sessionId;
/** Vars Context:业务变量 + 租户上下文,统一存储 */
private final Map<String, Object> varsContext = new ConcurrentHashMap<>();
/** Document Context:文件解析块列表,按 chunkIndex 顺序维护 */
private final List<DocumentChunk> documentChunks = new ArrayList<>();
/** 文档块去重索引:chunkId → 是否已存在 */
private final Map<String, Boolean> chunkIndex = new ConcurrentHashMap<>();
public InMemoryConversationMemory(String sessionId) {
this.sessionId = Objects.requireNonNull(sessionId, "sessionId 不能为 null");
}
@Override
public String getSessionId() {
return sessionId;
}
// ==================== Vars Context ====================
@Override
public void updateVariable(String key, Object value) {
if (key == null || key.isEmpty()) {
return;
}
if (value == null) {
varsContext.remove(key);
} else {
varsContext.put(key, value);
}
}
@Override
public Object getVariable(String key) {
if (key == null || key.isEmpty()) {
return null;
}
return varsContext.get(key);
}
@Override
@SuppressWarnings("unchecked")
public <T> T getVariable(String key, Class<T> targetType) {
Object value = getVariable(key);
if (value == null) {
return null;
}
if (targetType.isInstance(value)) {
return (T) value;
}
log.warn("[Memory] 变量类型不匹配 key={} expected={} actual={}",
key, targetType.getName(), value.getClass().getName());
return null;
}
@Override
public Map<String, Object> getVarsContext() {
// 返回不可变快照,外部修改需通过 updateVariable
return Collections.unmodifiableMap(new LinkedHashMap<>(varsContext));
}
@Override
public void bindTenantContext(TenantContext tenantContext) {
if (tenantContext == null) {
// 清除绑定的租户信息
varsContext.remove(KEY_TENANT_ID);
varsContext.remove(KEY_USER_ID);
varsContext.remove(KEY_ACCESS_TOKEN);
varsContext.remove(KEY_REQUEST_ATTRIBUTES);
varsContext.remove(KEY_SERVER_WEB_EXCHANGE);
return;
}
varsContext.put(KEY_TENANT_ID, tenantContext.getTenantId());
varsContext.put(KEY_USER_ID, tenantContext.getUserId());
varsContext.put(KEY_ACCESS_TOKEN, tenantContext.getAccessToken());
if (tenantContext.getRequestAttributes() != null) {
varsContext.put(KEY_REQUEST_ATTRIBUTES, tenantContext.getRequestAttributes());
}
if (tenantContext.getServerWebExchange() != null) {
varsContext.put(KEY_SERVER_WEB_EXCHANGE, tenantContext.getServerWebExchange());
}
}
@Override
public void clearVars() {
varsContext.clear();
}
// ==================== Document Context ====================
@Override
public void addDocumentChunks(List<DocumentChunk> chunks) {
if (chunks == null || chunks.isEmpty()) {
return;
}
synchronized (documentChunks) {
for (DocumentChunk chunk : chunks) {
if (chunk == null || chunk.getChunkId() == null) {
continue;
}
// chunkId 去重
if (chunkIndex.putIfAbsent(chunk.getChunkId(), Boolean.TRUE) != null) {
log.debug("[Memory] 文档块已存在,跳过 chunkId={} fileName={}",
chunk.getChunkId(), chunk.getFileName());
continue;
}
documentChunks.add(chunk);
}
// 按 chunkIndex 字段排序,保证顺序稳定
documentChunks.sort((a, b) -> Integer.compare(a.getChunkIndex(), b.getChunkIndex()));
}
}
@Override
public List<DocumentChunk> getDocumentChunks() {
synchronized (documentChunks) {
return new ArrayList<>(documentChunks);
}
}
@Override
public String getDocumentContext() {
List<DocumentChunk> snapshot;
synchronized (documentChunks) {
if (documentChunks.isEmpty()) {
return "";
}
snapshot = new ArrayList<>(documentChunks);
}
StringBuilder sb = new StringBuilder();
for (DocumentChunk chunk : snapshot) {
String fileName = chunk.getFileName() != null ? chunk.getFileName() : chunk.getFileKey();
sb.append("=== 文件: ").append(fileName).append(" ===\n");
sb.append(chunk.getContent() != null ? chunk.getContent() : "");
sb.append("\n=== 文件结束: ").append(fileName).append(" ===\n\n");
}
return sb.toString();
}
@Override
public void clearDocuments() {
synchronized (documentChunks) {
documentChunks.clear();
chunkIndex.clear();
}
}
// ==================== 上下文拼接 ====================
@Override
public String getSystemContext() {
StringBuilder sb = new StringBuilder();
// 拼接 Vars Context(仅展示 sys.tenant.* 和非内部变量)
Map<String, Object> vars = getVarsContext();
String tenantContext = buildTenantSection(vars);
if (!tenantContext.isEmpty()) {
sb.append("【用户上下文】\n").append(tenantContext).append("\n");
}
String bizVars = buildBizVarsSection(vars);
if (!bizVars.isEmpty()) {
sb.append("【业务变量】\n").append(bizVars).append("\n");
}
// 拼接 Document Context
String docContext = getDocumentContext();
if (!docContext.isEmpty()) {
sb.append("【上传文件】\n").append(docContext);
}
return sb.toString();
}
/**
* 构建租户上下文段落
*/
private String buildTenantSection(Map<String, Object> vars) {
Object tenantId = vars.get(KEY_TENANT_ID);
Object userId = vars.get(KEY_USER_ID);
if (tenantId == null && userId == null) {
return "";
}
StringBuilder sb = new StringBuilder();
if (tenantId != null) {
sb.append("租户: ").append(tenantId).append("\n");
}
if (userId != null) {
sb.append("用户: ").append(userId).append("\n");
}
return sb.toString();
}
/**
* 构建业务变量段落(过滤掉 sys.tenant.* 内部变量和 requestAttributes/exchange 等非文本对象)
*/
private String buildBizVarsSection(Map<String, Object> vars) {
StringBuilder sb = new StringBuilder();
vars.forEach((key, value) -> {
// 跳过租户上下文相关 key(已在用户上下文段落展示)
if (key.startsWith(TENANT_KEY_PREFIX)) {
return;
}
// 跳过非文本对象(RequestAttributes / ServerWebExchange 等)
if (value == null || value.getClass().getName().startsWith("org.springframework.web")) {
return;
}
// 跳过 Class 对象
if (value instanceof Class) {
return;
}
sb.append(key).append(": ").append(value).append("\n");
});
return sb.toString();
}
@Override
public void clearAll() {
// 仅清空本模块管理的两段记忆;Spring AI 的 ChatMemory 不在此清理
clearVars();
clearDocuments();
}
}
@@ -0,0 +1,75 @@
package cn.com.mfish.common.ai.memory;
import java.util.List;
import java.util.Objects;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.ConcurrentMap;
import java.util.stream.Collectors;
/**
* 基于 JVM 内存的 {@link ConversationMemoryStore} 默认实现
* <p>
* 使用 {@link ConcurrentHashMap} 按 sessionId 缓存 {@link ConversationMemory} 实例。
* 实例首次访问时创建并缓存,后续访问返回同一实例(sessionId 级单例)。
* </p>
* <p>
* 限制:进程内存储,重启丢失;不支持分布式部署下的会话漂移。
* 后续可替换为 Redis 实现以支持多实例共享。
* </p>
* <p>
* 注意:本 Store 只管理 Vars Context 和 Document Context,不管理 Spring AI 的
* ChatMemory(对话历史)。会话结束时如需清空对话历史,业务方需另行调用
* {@code ChatMemory.clear(sessionId)}。
* </p>
*
* @author: mfish
* @date: 2026/07/20
*/
public class InMemoryConversationMemoryStore implements ConversationMemoryStore {
private final ConcurrentMap<String, ConversationMemory> store = new ConcurrentHashMap<>();
public InMemoryConversationMemoryStore() {
}
@Override
public ConversationMemory getOrCreate(String sessionId) {
Objects.requireNonNull(sessionId, "sessionId 不能为 null");
if (sessionId.isEmpty()) {
throw new IllegalArgumentException("sessionId 不能为空字符串");
}
return store.computeIfAbsent(sessionId, InMemoryConversationMemory::new);
}
@Override
public ConversationMemory get(String sessionId) {
if (sessionId == null || sessionId.isEmpty()) {
return null;
}
return store.get(sessionId);
}
@Override
public ConversationMemory remove(String sessionId) {
if (sessionId == null || sessionId.isEmpty()) {
return null;
}
ConversationMemory removed = store.remove(sessionId);
if (removed != null) {
removed.clearAll();
}
return removed;
}
@Override
public List<String> listSessionIds() {
return store.keySet().stream()
.sorted()
.collect(Collectors.toList());
}
@Override
public int size() {
return store.size();
}
}
@@ -0,0 +1,33 @@
package cn.com.mfish.common.ai.memory;
import org.springframework.boot.autoconfigure.condition.ConditionalOnMissingBean;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
/**
* Memory 模块自动配置
* <p>
* 注册 {@link ConversationMemoryStore} 的默认实现 {@link InMemoryConversationMemoryStore}。
* </p>
* <p>
* 本模块只管理 Spring AI 管不到的两类上下文:Vars Context(变量/凭证)和
* Document Context(文件解析块)。对话历史(Short-term Memory)由 Spring AI 的
* {@code ChatMemory} + {@code MessageChatMemoryAdvisor} 自动管理,本模块不涉及。
* </p>
* <p>
* 业务方可通过自定义 {@link ConversationMemoryStore} Bean 覆盖默认实现,
* 例如替换为 Redis 版本以支持多实例部署。
* </p>
*
* @author: mfish
* @date: 2026/07/20
*/
@Configuration
public class MemoryAutoConfiguration {
@Bean
@ConditionalOnMissingBean
public ConversationMemoryStore conversationMemoryStore() {
return new InMemoryConversationMemoryStore();
}
}
@@ -1 +1,3 @@
cn.com.mfish.common.ai.engine.ApiToolAutoConfiguration
cn.com.mfish.common.ai.memory.MemoryAutoConfiguration
cn.com.mfish.common.ai.capability.CapabilityAutoConfiguration