feat: streamline workflow orchestration and chat routing

Add Zhipu web search integration, remove obsolete workflow nodes and resume handling, and separate model, agent, and workflow chat execution.
This commit is contained in:
ageerle
2026-07-30 00:15:44 +08:00
parent 83fd1ee983
commit 6e264ad500
45 changed files with 1217 additions and 1580 deletions

View File

@@ -24,6 +24,13 @@
<dependencies>
<!-- 智谱官方 Java SDK流程编排 Web Search 扩展节点 -->
<dependency>
<groupId>ai.z.openapi</groupId>
<artifactId>zai-sdk</artifactId>
<version>0.3.5</version>
</dependency>
<dependency>
<groupId>org.ruoyi</groupId>
<artifactId>ruoyi-common-chat</artifactId>
@@ -45,6 +52,11 @@
<artifactId>ruoyi-common-satoken</artifactId>
</dependency>
<dependency>
<groupId>org.ruoyi</groupId>
<artifactId>ruoyi-common-tenant</artifactId>
</dependency>
<dependency>
<groupId>org.ruoyi</groupId>
<artifactId>ruoyi-common-mail</artifactId>

View File

@@ -1,16 +1,13 @@
package org.ruoyi.workflow.controller;
import com.baomidou.mybatisplus.extension.plugins.pagination.Page;
import io.swagger.v3.oas.annotations.Operation;
import jakarta.annotation.Resource;
import jakarta.validation.constraints.Min;
import jakarta.validation.constraints.NotNull;
import org.ruoyi.common.core.domain.R;
import org.ruoyi.workflow.dto.workflow.WfRuntimeNodeDto;
import org.ruoyi.workflow.dto.workflow.WfRuntimeResp;
import org.ruoyi.workflow.dto.workflow.WorkflowResumeReq;
import org.ruoyi.workflow.service.WorkflowRuntimeService;
import org.ruoyi.workflow.workflow.WorkflowStarter;
import org.springframework.validation.annotation.Validated;
import org.springframework.web.bind.annotation.*;
@@ -24,16 +21,6 @@ public class WorkflowRuntimeController {
@Resource
private WorkflowRuntimeService workflowRuntimeService;
@Resource
private WorkflowStarter workflowStarter;
@Operation(summary = "接收用户输入以继续执行剩余流程")
@PostMapping(value = "/resume/{runtimeUuid}")
public R resume(@PathVariable String runtimeUuid, @RequestBody WorkflowResumeReq resumeReq) {
workflowStarter.resumeFlow(runtimeUuid, resumeReq.getFeedbackContent(), resumeReq.getSseEmitter());
return R.ok();
}
@GetMapping("/page")
public R<Page<WfRuntimeResp>> search(@RequestParam String wfUuid,
@NotNull @Min(1) Integer currentPage,

View File

@@ -337,7 +337,6 @@ public class AdiConstant {
public static final String DEFAULT_INPUT_PARAM_NAME = "input";
public static final String DEFAULT_OUTPUT_PARAM_NAME = "output";
public static final String DEFAULT_ERROR_OUTPUT_PARAM_NAME = "error_msg";
public static final String HUMAN_FEEDBACK_KEY = "human_feedback";
public static final int NODE_PROCESS_STATUS_READY = 1;
public static final int NODE_PROCESS_STATUS_DOING = 2;
public static final int NODE_PROCESS_STATUS_SUCCESS = 3;
@@ -347,7 +346,6 @@ public class AdiConstant {
public static final int WORKFLOW_PROCESS_STATUS_DOING = 2;
public static final int WORKFLOW_PROCESS_STATUS_SUCCESS = 3;
public static final int WORKFLOW_PROCESS_STATUS_FAIL = 4;
public static final int WORKFLOW_PROCESS_STATUS_WAITING_INPUT = 5;
public static final int WORKFLOW_NODE_PROCESS_TYPE_NORMAL = 1;
public static final int WORKFLOW_NODE_PROCESS_TYPE_CONDITIONAL = 2;

View File

@@ -1,10 +0,0 @@
package org.ruoyi.workflow.dto.workflow;
import lombok.Data;
import org.springframework.web.servlet.mvc.method.annotation.SseEmitter;
@Data
public class WorkflowResumeReq {
private String feedbackContent;
private SseEmitter sseEmitter;
}

View File

@@ -55,8 +55,13 @@ public class SSEEmitterHelper {
} else {
sseEmitter.send(msg);
}
} catch (IllegalStateException ise) {
// SSE连接已关闭用户刷新页面、关闭标签页或重新提交
log.warn("SSE emitter already completed for event [{}], ignoring", name);
COMPLETED_SSE.put(sseEmitter, Boolean.TRUE);
} catch (IOException ioException) {
log.error("stream onNext error", ioException);
COMPLETED_SSE.put(sseEmitter, Boolean.TRUE);
}
}

View File

@@ -4,7 +4,6 @@ import org.ruoyi.common.chat.domain.dto.request.ChatRequest;
import org.ruoyi.common.chat.domain.vo.chat.ChatModelVo;
import org.ruoyi.common.chat.enums.RoleType;
import lombok.extern.slf4j.Slf4j;
import org.ruoyi.common.core.exception.ServiceException;
import org.ruoyi.common.core.service.ConfigService;
import org.ruoyi.common.core.utils.SpringUtils;
import org.ruoyi.common.core.utils.StringUtils;
@@ -12,6 +11,7 @@ import org.ruoyi.workflow.entity.WorkflowNode;
import org.ruoyi.workflow.helper.SSEEmitterHelper;
import org.ruoyi.workflow.workflow.WfState;
import org.ruoyi.workflow.workflow.WorkflowUtil;
import org.ruoyi.workflow.workflow.node.enmus.NodeMessageTemplateEnum;
import org.springframework.web.servlet.mvc.method.annotation.SseEmitter;
/**
@@ -38,17 +38,20 @@ public class WorkflowMessageUtil {
/**
* 获取节点的响应模板
* 获取节点的响应模板 <br/>
* 优先读取 sys_config 配置(可在系统管理-配置管理中自定义),
* 未配置时回退到枚举内置默认模板, 模板仅为展示文案, 缺失不应中断工作流执行
* @param configKey 参数Key
* @return 返回模板样式
*/
public static String getNodeMessageTemplate(String configKey){
ConfigService configService = SpringUtil.getBean(ConfigService.class);
String configValue = configService.getConfigValue(configKey);
if (StringUtils.isEmpty(configValue)) {
throw new ServiceException("请先配置该节点的响应模板");
if (StringUtils.isNotEmpty(configValue)) {
return configValue;
}
return configValue;
log.warn("sys_config 未配置节点响应模板 [{}], 已回退使用内置默认模板", configKey);
return NodeMessageTemplateEnum.getDefaultTemplate(configKey);
}
/**

View File

@@ -1,17 +0,0 @@
package org.ruoyi.workflow.workflow;
import org.apache.commons.collections4.map.PassiveExpiringMap;
/**
* 已中断正在等待用户输入的流程 <br/>
* TODO 需要考虑项目多节点部署的情况
*/
public class InterruptedFlow {
/**
* 10分钟超时
*/
private static final PassiveExpiringMap.ExpirationPolicy<String, WorkflowEngine> ep = new PassiveExpiringMap.ConstantTimeToLiveExpirationPolicy<>(60 * 1000 * 10);
public static PassiveExpiringMap<String, WorkflowEngine> RUNTIME_TO_GRAPH = new PassiveExpiringMap<>(ep);
}

View File

@@ -16,24 +16,14 @@ public enum WfComponentNameEnum {
TONGYI_WANX("Tongyiwanx"),
DOCUMENT_EXTRACTOR("DocumentExtractor"),
KEYWORD_EXTRACTOR("KeywordExtractor"),
FAQ_EXTRACTOR("FaqExtractor"),
KNOWLEDGE_RETRIEVER("KnowledgeRetrieval"),
SWITCHER("Switcher"),
CLASSIFIER("Classifier"),
TEMPLATE("Template"),
GOOGLE_SEARCH("Google"),
HUMAN_FEEDBACK("HumanFeedback"),
MAIL_SEND("MailSend"),
HTTP_REQUEST("HttpRequest");

View File

@@ -4,15 +4,14 @@ import org.ruoyi.workflow.entity.WorkflowComponent;
import org.ruoyi.workflow.entity.WorkflowNode;
import org.ruoyi.workflow.workflow.node.AbstractWfNode;
import org.ruoyi.workflow.workflow.node.EndNode;
import org.ruoyi.workflow.workflow.node.humanFeedBack.HumanFeedbackNode;
import org.ruoyi.workflow.workflow.node.answer.LLMAnswerNode;
import org.ruoyi.workflow.workflow.node.httpRequest.HttpRequestNode;
import org.ruoyi.workflow.workflow.node.image.ImageNode;
import org.ruoyi.workflow.workflow.node.keywordExtractor.KeywordExtractorNode;
import org.ruoyi.workflow.workflow.node.knowledgeRetrieval.KnowledgeRetrievalNode;
import org.ruoyi.workflow.workflow.node.mailSend.MailSendNode;
import org.ruoyi.workflow.workflow.node.start.StartNode;
import org.ruoyi.workflow.workflow.node.switcher.SwitcherNode;
import org.ruoyi.workflow.workflow.node.googleSearch.GoogleSearchNode;
public class WfNodeFactory {
public static AbstractWfNode create(WorkflowComponent wfComponent, WorkflowNode nodeDefinition,
@@ -21,14 +20,13 @@ public class WfNodeFactory {
switch (WfComponentNameEnum.getByName(wfComponent.getName())) {
case START -> wfNode = new StartNode(wfComponent, nodeDefinition, wfState, nodeState);
case LLM_ANSWER -> wfNode = new LLMAnswerNode(wfComponent, nodeDefinition, wfState, nodeState);
case KEYWORD_EXTRACTOR -> wfNode = new KeywordExtractorNode(wfComponent, nodeDefinition, wfState, nodeState);
case TONGYI_WANX -> wfNode = new ImageNode(wfComponent, nodeDefinition, wfState, nodeState);
case KNOWLEDGE_RETRIEVER -> wfNode = new KnowledgeRetrievalNode(wfComponent, nodeDefinition, wfState, nodeState);
case END -> wfNode = new EndNode(wfComponent, nodeDefinition, wfState, nodeState);
case MAIL_SEND -> wfNode = new MailSendNode(wfComponent, nodeDefinition, wfState, nodeState);
case HTTP_REQUEST -> wfNode = new HttpRequestNode(wfComponent, nodeDefinition, wfState, nodeState);
case SWITCHER -> wfNode = new SwitcherNode(wfComponent, nodeDefinition, wfState, nodeState);
case HUMAN_FEEDBACK -> wfNode = new HumanFeedbackNode(wfComponent, nodeDefinition, wfState, nodeState);
case GOOGLE_SEARCH -> wfNode = new GoogleSearchNode(wfComponent, nodeDefinition, wfState, nodeState);
default -> {
}
}

View File

@@ -55,11 +55,6 @@ public class WfState {
private List<NodeIOData> output = new ArrayList<>();
private Integer processStatus = WORKFLOW_PROCESS_STATUS_READY;
/**
* 人机交互节点
*/
private Set<String> interruptNodes = new HashSet<>();
public WfState(User user, List<NodeIOData> input, String uuid, Long userId, String tokenValue, SseEmitter sseEmitter, Long sessionId) {
this.input = input;
this.user = user;
@@ -133,8 +128,4 @@ public class WfState {
.findFirst()
.orElse(null);
}
public void addInterruptNode(String nodeUuid) {
this.interruptNodes.add(nodeUuid);
}
}

View File

@@ -108,82 +108,39 @@ public class WorkflowEngine {
MemorySaver saver = new MemorySaver();
CompileConfig compileConfig = CompileConfig.builder().checkpointSaver(saver)
.interruptBefore(wfState.getInterruptNodes().toArray(String[]::new))
.build();
app = mainStateGraph.compile(compileConfig);
RunnableConfig invokeConfig = RunnableConfig.builder().build();
exe(invokeConfig, false);
exe(invokeConfig);
} catch (Exception e) {
errorWhenExe(e);
}
}
private void exe(RunnableConfig invokeConfig, boolean resume) {
private void exe(RunnableConfig invokeConfig) {
//不使用langgraph4j state的update相关方法无需传入input
AsyncGenerator<NodeOutput<WfNodeState>> outputs = app.stream(resume ? null : Map.of(), invokeConfig);
AsyncGenerator<NodeOutput<WfNodeState>> outputs = app.stream(Map.of(), invokeConfig);
streamingResult(wfState, outputs, sseEmitter);
StateSnapshot<WfNodeState> stateSnapshot = app.getState(invokeConfig);
String nextNode = stateSnapshot.config().nextNode().orElse("");
//还有下个节点表示进入中断状态等待用户输入后继续执<E7BBAD>?
if (StringUtils.isNotBlank(nextNode) && !nextNode.equalsIgnoreCase(END)) {
// 获取提示模板
String nodeMessageTemplate = WorkflowMessageUtil.getNodeMessageTemplate(NodeMessageTemplateEnum.HUMAN_FEED_BACK.getValue());
// 获取人机交互提示信息
String intTip = nodeMessageTemplate + WorkflowUtil.getHumanFeedbackTip(nextNode, wfNodes);
//将等待输入信息[事件与提示词]发送到到客户端
SSEEmitterHelper.parseAndSendPartialMsg(sseEmitter, "[NODE_WAIT_FEEDBACK_BY_" + nextNode + "]", intTip);
// 保存提示信息到Chat信息记录中对话使用
WorkflowMessageUtil.saveWorkflowMessage(wfState, intTip);
InterruptedFlow.RUNTIME_TO_GRAPH.put(wfState.getUuid(), this);
//更新状<E696B0>?
wfState.setProcessStatus(WORKFLOW_PROCESS_STATUS_WAITING_INPUT);
workflowRuntimeService.updateOutput(wfRuntimeResp.getId(), wfState);
} else {
WorkflowRuntime updatedRuntime = workflowRuntimeService.updateOutput(wfRuntimeResp.getId(), wfState);
// 保存成功会话信息
wfNodes.stream().filter(item -> stateSnapshot.node().equals(item.getUuid()))
.findFirst().ifPresent(wfNode -> {
// 获取节点模板提示词信息
String nodeMessageTemplate = WorkflowMessageUtil.getNodeMessageTemplate(NodeMessageTemplateEnum.END.getValue());
// 发送SSE消息驱动事件和保存会话
WorkflowMessageUtil.notifyAndStoreMessage(wfState, sseEmitter, wfNode, nodeMessageTemplate);
});
// 发送结束消息
sseEmitterHelper.sendComplete(user.getId(), sseEmitter, updatedRuntime.getOutput());
// 发送驱动消息事件
InterruptedFlow.RUNTIME_TO_GRAPH.remove(wfState.getUuid());
}
}
/**
* 中断流程等待用户输入时,会进行暂停状态,用户输入后调用本方法执行流程剩余部分
*
* @param userInput 用户输入
*/
public void resume(String userInput) {
RunnableConfig invokeConfig = RunnableConfig.builder().build();
try {
app.updateState(invokeConfig, Map.of(HUMAN_FEEDBACK_KEY, userInput), null);
exe(invokeConfig, true);
} catch (Exception e) {
errorWhenExe(e);
} finally {
//有可能多次接收人机交互,待整个流程完全执行后才能删除
if (wfState.getProcessStatus() != WORKFLOW_PROCESS_STATUS_WAITING_INPUT) {
InterruptedFlow.RUNTIME_TO_GRAPH.remove(wfState.getUuid());
}
}
wfState.setProcessStatus(WORKFLOW_PROCESS_STATUS_SUCCESS);
WorkflowRuntime updatedRuntime = workflowRuntimeService.updateOutput(wfRuntimeResp.getId(), wfState);
wfNodes.stream().filter(item -> stateSnapshot.node().equals(item.getUuid()))
.findFirst().ifPresent(wfNode -> {
String nodeMessageTemplate = WorkflowMessageUtil.getNodeMessageTemplate(NodeMessageTemplateEnum.END.getValue());
WorkflowMessageUtil.notifyAndStoreMessage(wfState, sseEmitter, wfNode, nodeMessageTemplate);
});
sseEmitterHelper.sendComplete(user.getId(), sseEmitter, updatedRuntime.getOutput());
}
private void errorWhenExe(Exception e) {
log.error("error", e);
String nodeMessageTemplate = WorkflowMessageUtil.getNodeMessageTemplate(NodeMessageTemplateEnum.EXCEPTION.getValue());
String errorMsg = e.getMessage();
if (errorMsg.contains("parallel node doesn't support conditional branch")) {
if (errorMsg != null && errorMsg.contains("parallel node doesn't support conditional branch")) {
errorMsg = "并行节点中不能包含条件分<EFBFBD>?";
}
errorMsg = nodeMessageTemplate + errorMsg;
errorMsg = nodeMessageTemplate + (errorMsg != null ? errorMsg : e.getClass().getSimpleName());
// 保存会话信息且发送驱动消息事件
WorkflowMessageUtil.saveWorkflowMessage(wfState, errorMsg);
sseEmitterHelper.sendErrorAndComplete(user.getId(), sseEmitter, errorMsg);
@@ -267,6 +224,10 @@ public class WorkflowEngine {
log.info("node:{},chunk:{}", node, chunk);
SSEEmitterHelper.parseAndSendPartialMsg(sseEmitter, "[NODE_CHUNK_" + node + "]", chunk);
} else {
// __END__ 是 langgraph4j 的终止伪节点, 无对应业务节点状态, 跳过
if (END.equals(out.node())) {
continue;
}
AbstractWfNode abstractWfNode = wfState.getCompletedNodes().stream()
.filter(item -> item.getNode().getUuid().endsWith(out.node())).findFirst().orElse(null);
if (null != abstractWfNode) {

View File

@@ -19,7 +19,6 @@ import static org.bsc.langgraph4j.StateGraph.END;
import static org.bsc.langgraph4j.StateGraph.START;
import static org.bsc.langgraph4j.action.AsyncEdgeAction.edge_async;
import static org.bsc.langgraph4j.action.AsyncNodeAction.node_async;
import static org.ruoyi.workflow.workflow.WfComponentNameEnum.HUMAN_FEEDBACK;
/**
* 负责构建工作流运行所依赖的状态图<E68081>?
@@ -27,7 +26,6 @@ import static org.ruoyi.workflow.workflow.WfComponentNameEnum.HUMAN_FEEDBACK;
@Slf4j
public class WorkflowGraphBuilder {
private final Map<Long, WorkflowComponent> componentIndex;
private final Map<String, WorkflowNode> nodeIndex;
private final Map<String, List<WorkflowEdge>> edgesBySource;
private final Map<String, List<WorkflowEdge>> edgesByTarget;
@@ -46,8 +44,6 @@ public class WorkflowGraphBuilder {
List<WorkflowEdge> edges,
WorkflowNodeRunner nodeRunner,
WfState wfState) {
this.componentIndex = components.stream()
.collect(Collectors.toMap(WorkflowComponent::getId, Function.identity(), (origin, ignore) -> origin));
this.nodeIndex = nodes.stream()
.collect(Collectors.toMap(WorkflowNode::getUuid, Function.identity(), (origin, ignore) -> origin));
this.edgesBySource = edges.stream().collect(Collectors.groupingBy(WorkflowEdge::getSourceNodeUuid));
@@ -217,14 +213,6 @@ public class WorkflowGraphBuilder {
WorkflowNode wfNode = getNodeByUuid(stateGraphNodeUuid);
stateGraph.addNode(stateGraphNodeUuid, node_async(state -> nodeRunner.run(wfNode, state)));
stateGraphList.add(stateGraph);
WorkflowComponent component = componentIndex.get(wfNode.getWorkflowComponentId());
if (component == null) {
throw new BaseException(ErrorEnum.A_PARAMS_ERROR.getInfo());
}
if (HUMAN_FEEDBACK.getName().equals(component.getName())) {
wfState.addInterruptNode(stateGraphNodeUuid);
}
}
private void addEdgeToStateGraph(StateGraph<WfNodeState> stateGraph, String source, String target) throws GraphStateException {

View File

@@ -6,9 +6,9 @@ import jakarta.annotation.Resource;
import lombok.extern.slf4j.Slf4j;
import org.ruoyi.common.chat.entity.User;
import org.ruoyi.common.chat.service.workFlow.IWorkFlowStarterService;
import org.ruoyi.common.core.exception.base.BaseException;
import org.ruoyi.common.satoken.utils.LoginHelper;
import org.ruoyi.common.sse.core.SseEmitterManager;
import org.ruoyi.common.tenant.helper.TenantHelper;
import org.ruoyi.workflow.entity.*;
import org.ruoyi.workflow.helper.SSEEmitterHelper;
import org.ruoyi.workflow.service.*;
@@ -58,6 +58,8 @@ public class WorkflowStarter implements IWorkFlowStarterService {
Long userId = LoginHelper.getUserId();
// 获取登录Token仅透传给 WfState工作流 SSE 通过 emitter 直发,不串台)
String tokenValue = StpUtil.getTokenValue();
// 获取当前租户ID@Async 线程不继承请求线程的租户上下文,需显式透传)
String tenantId = TenantHelper.getTenantId();
// 根据会话ID连接SSE对象每会话一个连接避免同用户多会话串台
SseEmitter sseEmitter = sseEmitterManager.connect(String.valueOf(sessionId));
if (!sseEmitterHelper.checkOrComplete(user, sseEmitter)) {
@@ -71,41 +73,33 @@ public class WorkflowStarter implements IWorkFlowStarterService {
sseEmitterHelper.sendErrorAndComplete(user.getId(), sseEmitter, A_WF_DISABLED.getInfo());
return sseEmitter;
}
self.asyncRun(user, workflow, userInputs, sseEmitter, userId, tokenValue, sessionId);
self.asyncRun(user, workflow, userInputs, sseEmitter, userId, tokenValue, sessionId, tenantId);
return sseEmitter;
}
@Async
public void asyncRun(User user, Workflow workflow, List<ObjectNode> userInputs, SseEmitter sseEmitter, Long userId, String tokenValue, Long sessionId) {
log.info("WorkflowEngine run,userId:{},workflowUuid:{},userInputs:{}", user.getId(), workflow.getUuid(), userInputs);
List<WorkflowComponent> components = workflowComponentService.getAllEnable();
List<WorkflowNode> nodes = workflowNodeService.lambdaQuery()
.eq(WorkflowNode::getWorkflowId, workflow.getId())
.eq(WorkflowNode::getIsDeleted, false)
.list();
List<WorkflowEdge> edges = workflowEdgeService.lambdaQuery()
.eq(WorkflowEdge::getWorkflowId, workflow.getId())
.eq(WorkflowEdge::getIsDeleted, false)
.list();
WorkflowEngine workflowEngine = new WorkflowEngine(workflow,
sseEmitterHelper, components, nodes, edges,
workflowRuntimeService, workflowRuntimeNodeService);
workflowEngine.run(user, userInputs, sseEmitter, userId, tokenValue, sessionId);
}
@Async
public void resumeFlow(String runtimeUuid, String userInput, SseEmitter sseEmitter) {
WorkflowEngine workflowEngine = InterruptedFlow.RUNTIME_TO_GRAPH.get(runtimeUuid);
if (null == workflowEngine) {
log.error("工作流恢复执行时失败,runtime:{}", runtimeUuid);
throw new BaseException(A_WF_RESUME_FAIL.getInfo());
public void asyncRun(User user, Workflow workflow, List<ObjectNode> userInputs, SseEmitter sseEmitter, Long userId, String tokenValue, Long sessionId, String tenantId) {
// @Async 线程不继承请求线程的租户上下文, 显式设置, 避免租户缓存/隔离逻辑异常
if (tenantId != null) {
TenantHelper.setDynamic(tenantId);
}
// 如果SSE连接对象不为空传入该对象Chat调用工作流对话使用
if (null != sseEmitter){
workflowEngine.setSseEmitter(sseEmitter);
// 为了让每个节点都可以发送模板消息 保持SSE对象一致以防出现向已关闭的SSE对象发送消息
workflowEngine.getWfState().setSseEmitter(sseEmitter);
try {
log.info("WorkflowEngine run,userId:{},workflowUuid:{},userInputs:{}", user.getId(), workflow.getUuid(), userInputs);
List<WorkflowComponent> components = workflowComponentService.getAllEnable();
List<WorkflowNode> nodes = workflowNodeService.lambdaQuery()
.eq(WorkflowNode::getWorkflowId, workflow.getId())
.eq(WorkflowNode::getIsDeleted, false)
.list();
List<WorkflowEdge> edges = workflowEdgeService.lambdaQuery()
.eq(WorkflowEdge::getWorkflowId, workflow.getId())
.eq(WorkflowEdge::getIsDeleted, false)
.list();
WorkflowEngine workflowEngine = new WorkflowEngine(workflow,
sseEmitterHelper, components, nodes, edges,
workflowRuntimeService, workflowRuntimeNodeService);
workflowEngine.run(user, userInputs, sseEmitter, userId, tokenValue, sessionId);
} finally {
TenantHelper.clearDynamic();
}
workflowEngine.resume(userInput);
}
}

View File

@@ -2,7 +2,6 @@ package org.ruoyi.workflow.workflow;
import cn.hutool.core.collection.CollStreamUtil;
import cn.hutool.core.collection.CollUtil;
import cn.hutool.core.util.StrUtil;
import dev.langchain4j.data.message.ChatMessage;
import dev.langchain4j.data.message.UserMessage;
import dev.langchain4j.model.chat.response.StreamingChatResponseHandler;
@@ -20,7 +19,6 @@ import org.ruoyi.common.chat.factory.ImageServiceFactory;
import org.ruoyi.workflow.base.NodeInputConfigTypeHandler;
import org.ruoyi.workflow.entity.WorkflowNode;
import org.ruoyi.workflow.enums.WfIODataTypeEnum;
import org.ruoyi.workflow.util.JsonUtil;
import org.ruoyi.workflow.workflow.data.NodeIOData;
import org.ruoyi.workflow.workflow.data.NodeIODataContent;
import org.ruoyi.workflow.workflow.def.WfNodeParamRef;
@@ -88,22 +86,6 @@ public class WorkflowUtil{
return result;
}
public static String getHumanFeedbackTip(String nodeUuid, List<WorkflowNode> wfNodes) {
WorkflowNode wfNode = wfNodes.stream()
.filter(item -> item.getUuid().equals(nodeUuid))
.findFirst().orElse(null);
if (null == wfNode) {
return "";
}
String wfNodeNodeConfig = wfNode.getNodeConfig();
if (StrUtil.isBlank(wfNodeNodeConfig)) {
return "";
}
Map<String, Object> map = JsonUtil.toMap(wfNodeNodeConfig);
Object tip = map.getOrDefault("tip", "");
return String.valueOf(tip);
}
public void streamingInvokeLLM(WfState wfState, WfNodeState state, WorkflowNode node, String modelName,
String prompt, String nodeMessageTemplate) {
log.info("stream invoke, modelName: {}", modelName);

View File

@@ -1,16 +0,0 @@
package org.ruoyi.workflow.workflow.node.classifier;
import com.fasterxml.jackson.annotation.JsonProperty;
import lombok.Data;
import java.util.ArrayList;
import java.util.List;
@Data
public class ClassifierNodeConfig {
private List categories = new ArrayList<>();
@JsonProperty("model_platform")
private String modelPlatform;
@JsonProperty("model_name")
private String modelName;
}

View File

@@ -3,23 +3,44 @@ package org.ruoyi.workflow.workflow.node.enmus;
import lombok.Getter;
/**
* 节点消息模板ConfigKey枚举
* 节点消息模板ConfigKey枚举 <br/>
* 模板优先从 sys_config 读取(可在系统管理-配置管理中自定义), 未配置时回退到 defaultTemplate 内置默认值
*/
@Getter
public enum NodeMessageTemplateEnum {
HTTP_REQUEST("node.httpRequest.template"),
MAIL_SEND("node.mailsend.template"),
IMAGE("node.image.template"),
HUMAN_FEED_BACK("node.humanFeedback.template"),
SWITCH("node.switch.template"),
LLM_RESPONSE("node.llmAnswer.template"),
KEYWORD_EXTRACTOR("node.keywordExtractor.template"),
EXCEPTION("node.exception.template"),
END("node.end.template");
HTTP_REQUEST("node.httpRequest.template", "✅ HTTP请求节点结束响应 - "),
MAIL_SEND("node.mailsend.template", "📧 发送邮箱节点:结束响应 - "),
IMAGE("node.image.template", "🎨 文生图节点:结束响应 - 图片URL: "),
SWITCH("node.switch.template", "🔀 条件分支节点:触发 -> 跳转到节点 "),
LLM_RESPONSE("node.llmAnswer.template", "🤖 LLM 节点 生成回答:"),
GOOGLE_SEARCH("node.googleSearch.template", "🔍 网络搜索节点处理完成:"),
EXCEPTION("node.exception.template", "🛑 工作流发生异常:"),
END("node.end.template", "🔚 流程已执行完毕,如果您有其他需求,请随时重新发起请求。");
private final String value;
NodeMessageTemplateEnum(String value) {
/**
* 内置默认模板, sys_config 未配置对应键时使用
*/
private final String defaultTemplate;
NodeMessageTemplateEnum(String value, String defaultTemplate) {
this.value = value;
this.defaultTemplate = defaultTemplate;
}
/**
* 根据 configKey 获取内置默认模板, 未知键返回空串
*
* @param configKey sys_config 配置键
* @return 内置默认模板
*/
public static String getDefaultTemplate(String configKey) {
for (NodeMessageTemplateEnum item : values()) {
if (item.value.equals(configKey)) {
return item.defaultTemplate;
}
}
return "";
}
}

View File

@@ -0,0 +1,91 @@
package org.ruoyi.workflow.workflow.node.googleSearch;
import lombok.extern.slf4j.Slf4j;
import org.apache.commons.lang3.StringUtils;
import org.ruoyi.workflow.entity.WorkflowComponent;
import org.ruoyi.workflow.entity.WorkflowNode;
import org.ruoyi.workflow.util.JsonUtil;
import org.ruoyi.workflow.util.SpringUtil;
import org.ruoyi.workflow.workflow.NodeProcessResult;
import org.ruoyi.workflow.workflow.WfNodeState;
import org.ruoyi.workflow.workflow.WfState;
import org.ruoyi.workflow.workflow.WorkflowUtil;
import org.ruoyi.workflow.workflow.data.NodeIOData;
import org.ruoyi.workflow.workflow.node.AbstractWfNode;
import org.ruoyi.workflow.workflow.node.enmus.NodeMessageTemplateEnum;
import java.util.List;
import java.util.UUID;
import static org.ruoyi.workflow.cosntant.AdiConstant.WorkflowConstant.DEFAULT_OUTPUT_PARAM_NAME;
/**
* 【扩展节点】网络搜索
* 通过智谱 Web Search API 返回适合大模型消费的结构化网页结果。
*/
@Slf4j
public class GoogleSearchNode extends AbstractWfNode {
public GoogleSearchNode(WorkflowComponent wfComponent, WorkflowNode nodeDef, WfState wfState, WfNodeState nodeState) {
super(wfComponent, nodeDef, wfState, nodeState);
}
/**
* 处理搜索请求
* nodeConfig 格式:
* {
* "query": "搜索关键词",
* "search_engine": "search_std",
* "result_count": 10,
* "search_domain_filter": "",
* "search_recency_filter": "noLimit",
* "content_size": "medium",
* "include_image": false
* }
*
* @return 搜索结果
*/
@Override
public NodeProcessResult onProcess() {
GoogleSearchNodeConfig config = checkAndGetConfig(GoogleSearchNodeConfig.class);
// 获取搜索关键词
String searchQuery = WorkflowUtil.renderTemplate(config.getQuery(), state.getInputs());
if (StringUtils.isBlank(searchQuery)) {
searchQuery = getFirstInputText();
}
if (StringUtils.isBlank(searchQuery)) {
throw new IllegalArgumentException("未提供搜索关键词");
}
searchQuery = searchQuery.trim();
if (searchQuery.length() > 70) {
throw new IllegalArgumentException("搜索关键词不能超过 70 个字符");
}
log.info("Web search node processing, engine: {}, result_count: {}",
config.getSearchEngine(), config.getResultCount());
String nodeMessageTemplate = getNodeMessageTemplate(NodeMessageTemplateEnum.GOOGLE_SEARCH.getValue());
notifyAndStoreMessage(wfState, nodeMessageTemplate);
ZhipuWebSearchClient searchClient = SpringUtil.getBean(ZhipuWebSearchClient.class);
ZhipuWebSearchClient.SearchResponse response = searchClient.search(
searchQuery,
config,
UUID.randomUUID().toString()
);
String searchResult = JsonUtil.toJson(response);
if (searchResult == null) {
throw new IllegalStateException("搜索结果序列化失败");
}
log.info("Web search completed, result count: {}", response.count());
notifyAndStoreMessage(wfState, nodeMessageTemplate + "返回 " + response.count() + " 条结果");
List<NodeIOData> outputs = List.of(
NodeIOData.createByText(DEFAULT_OUTPUT_PARAM_NAME, "智谱网络搜索结果", searchResult)
);
return NodeProcessResult.builder().content(outputs).build();
}
}

View File

@@ -0,0 +1,44 @@
package org.ruoyi.workflow.workflow.node.googleSearch;
import com.fasterxml.jackson.annotation.JsonProperty;
import jakarta.validation.constraints.Max;
import jakarta.validation.constraints.Min;
import jakarta.validation.constraints.Pattern;
import lombok.Data;
@Data
public class GoogleSearchNodeConfig {
/**
* 搜索查询关键词
*/
private String query;
@JsonProperty("search_engine")
@Pattern(
regexp = "search_std|search_pro|search_pro_sogou|search_pro_quark",
message = "搜索引擎参数无效"
)
private String searchEngine = "search_std";
@JsonProperty("result_count")
@Min(value = 1, message = "搜索结果数量不能小于 1")
@Max(value = 50, message = "搜索结果数量不能大于 50")
private Integer resultCount = 10;
@JsonProperty("search_domain_filter")
private String searchDomainFilter;
@JsonProperty("search_recency_filter")
@Pattern(
regexp = "oneDay|oneWeek|oneMonth|oneYear|noLimit",
message = "搜索时间范围参数无效"
)
private String searchRecencyFilter = "noLimit";
@JsonProperty("content_size")
@Pattern(regexp = "medium|high", message = "网页摘要长度参数无效")
private String contentSize = "medium";
@JsonProperty("include_image")
private Boolean includeImage = false;
}

View File

@@ -0,0 +1,155 @@
package org.ruoyi.workflow.workflow.node.googleSearch;
import ai.z.openapi.ZhipuAiClient;
import ai.z.openapi.service.web_search.WebSearchRequest;
import ai.z.openapi.service.web_search.WebSearchResp;
import ai.z.openapi.service.web_search.WebSearchResponse;
import lombok.RequiredArgsConstructor;
import org.apache.commons.lang3.StringUtils;
import org.ruoyi.common.chat.domain.bo.chat.ChatModelBo;
import org.ruoyi.common.chat.service.chat.IChatModelService;
import org.springframework.stereotype.Component;
import java.util.List;
import java.util.concurrent.TimeUnit;
/**
* 智谱 Web Search 官方 SDK 适配器。
*/
@Component
@RequiredArgsConstructor
public class ZhipuWebSearchClient {
private static final String ZHIPU_PROVIDER_CODE = "zhipu";
private static final String DEFAULT_BASE_URL = "https://open.bigmodel.cn/api/paas/v4/";
private final ZhipuWebSearchProperties properties;
private final IChatModelService chatModelService;
public SearchResponse search(String query, GoogleSearchNodeConfig config, String requestId) {
Credential credential = resolveCredential();
ZhipuAiClient client = createClient(credential);
try {
WebSearchRequest request = WebSearchRequest.builder()
.searchQuery(query)
.searchEngine(config.getSearchEngine())
.count(config.getResultCount())
.searchDomainFilter(blankToNull(config.getSearchDomainFilter()))
.searchRecencyFilter(config.getSearchRecencyFilter())
.contentSize(config.getContentSize())
.includeImage(config.getIncludeImage())
.requestId(requestId)
.build();
WebSearchResponse response = client.webSearch().createWebSearch(request);
if (response == null || !response.isSuccess() || response.getData() == null) {
String message = response == null ? "接口未返回响应" : StringUtils.defaultIfBlank(response.getMsg(), "未知错误");
throw new IllegalStateException("智谱 Web Search 调用失败:" + message);
}
List<SearchResult> results = response.getData().getWebSearchResp() == null
? List.of()
: response.getData().getWebSearchResp().stream()
.map(this::toSearchResult)
.toList();
return new SearchResponse(
query,
config.getSearchEngine(),
response.getData().getRequestId(),
results.size(),
results
);
} finally {
client.close();
}
}
private ZhipuAiClient createClient(Credential credential) {
int connectTimeout = positiveOrDefault(properties.getConnectTimeout(), 10);
int readTimeout = positiveOrDefault(properties.getReadTimeout(), 30);
return ZhipuAiClient.builder()
.ofZHIPU()
.apiKey(credential.apiKey())
.baseUrl(credential.baseUrl())
.networkConfig(connectTimeout, readTimeout, readTimeout, readTimeout, TimeUnit.SECONDS)
.enableTokenCache()
.build();
}
private Credential resolveCredential() {
if (isUsableApiKey(properties.getApiKey())) {
return new Credential(normalizeBaseUrl(properties.getBaseUrl()), properties.getApiKey().trim());
}
ChatModelBo query = new ChatModelBo();
query.setProviderCode(ZHIPU_PROVIDER_CODE);
return chatModelService.queryList(query).stream()
.filter(model -> isUsableApiKey(model.getApiKey()))
.findFirst()
.map(model -> new Credential(normalizeBaseUrl(model.getApiHost()), model.getApiKey().trim()))
.orElseThrow(() -> new IllegalStateException(
"未配置智谱 API Key请设置环境变量 ZAI_API_KEY或在模型管理中配置 zhipu 厂商密钥"
));
}
private boolean isUsableApiKey(String apiKey) {
return StringUtils.isNotBlank(apiKey)
&& !"sk_xx".equalsIgnoreCase(apiKey.trim())
&& !"your_api_key".equalsIgnoreCase(apiKey.trim());
}
private String normalizeBaseUrl(String baseUrl) {
String normalized = StringUtils.defaultIfBlank(baseUrl, DEFAULT_BASE_URL).trim();
normalized = StringUtils.removeEnd(normalized, "/");
if (!normalized.endsWith("/api/paas/v4")) {
normalized += "/api/paas/v4";
}
return normalized + "/";
}
private int positiveOrDefault(Integer value, int defaultValue) {
return value != null && value > 0 ? value : defaultValue;
}
private String blankToNull(String value) {
return StringUtils.isBlank(value) ? null : value.trim();
}
private SearchResult toSearchResult(WebSearchResp result) {
return new SearchResult(
result.getTitle(),
result.getContent(),
result.getLink(),
result.getMedia(),
result.getIcon(),
result.getRefer(),
result.getPublishDate(),
result.getImages()
);
}
private record Credential(String baseUrl, String apiKey) {
}
public record SearchResponse(
String query,
String searchEngine,
String requestId,
int count,
List<SearchResult> results
) {
}
public record SearchResult(
String title,
String content,
String link,
String media,
String icon,
String refer,
String publishDate,
List<String> images
) {
}
}

View File

@@ -0,0 +1,34 @@
package org.ruoyi.workflow.workflow.node.googleSearch;
import lombok.Data;
import org.springframework.boot.context.properties.ConfigurationProperties;
import org.springframework.stereotype.Component;
/**
* 智谱 Web Search 配置。
*/
@Data
@Component
@ConfigurationProperties(prefix = "workflow.web-search.zhipu")
public class ZhipuWebSearchProperties {
/**
* 智谱国内开放平台 API 根地址。
*/
private String baseUrl = "https://open.bigmodel.cn/api/paas/v4/";
/**
* 智谱 API Key。建议通过环境变量 ZAI_API_KEY 注入。
*/
private String apiKey;
/**
* 连接超时秒数。
*/
private Integer connectTimeout = 10;
/**
* 读取超时秒数。
*/
private Integer readTimeout = 30;
}

View File

@@ -1,56 +0,0 @@
package org.ruoyi.workflow.workflow.node.humanFeedBack;
import lombok.extern.slf4j.Slf4j;
import org.apache.commons.lang3.StringUtils;
import org.ruoyi.workflow.entity.WorkflowComponent;
import org.ruoyi.workflow.entity.WorkflowNode;
import org.ruoyi.workflow.workflow.NodeProcessResult;
import org.ruoyi.workflow.workflow.WfNodeState;
import org.ruoyi.workflow.workflow.WfState;
import org.ruoyi.workflow.workflow.WorkflowUtil;
import org.ruoyi.workflow.workflow.data.NodeIOData;
import org.ruoyi.workflow.workflow.node.AbstractWfNode;
import static org.ruoyi.workflow.cosntant.AdiConstant.WorkflowConstant.*;
/**
* 人机交互节点实现类
*/
@Slf4j
public class HumanFeedbackNode extends AbstractWfNode {
public HumanFeedbackNode(WorkflowComponent component, WorkflowNode nodeDefinition, WfState wfState, WfNodeState nodeState) {
super(component, nodeDefinition, wfState, nodeState);
}
// 人机交互节点的处理逻辑
@Override
public NodeProcessResult onProcess() {
log.info("Processing HumanFeedback node: {}", node.getTitle());
// 从状态中获取用户输入数据
Object humanFeedbackState = state.data().get(HUMAN_FEEDBACK_KEY);
if (null != humanFeedbackState) {
String userInput = humanFeedbackState.toString();
if (StringUtils.isNotBlank(userInput)) {
// 用户已提供输入,将用户输入添加到节点输入和输出中
NodeIOData feedbackData = NodeIOData.createByText("output", "default", userInput);
// 添加到输出列表,这样后续节点可以使用
state.getOutputs().add(feedbackData);
// 设置为成功状态
state.setProcessStatus(NODE_PROCESS_STATUS_SUCCESS);
log.info("Human feedback processed for node: {}, content: {}", node.getTitle(), userInput);
} else {
// 用户输入为空,设置等待状态
state.setProcessStatus(NODE_PROCESS_STATUS_DOING);
log.info("Human feedback is empty for node: {}", node.getTitle());
}
} else {
// 没有用户输入,这可能是正常情况(等待用户输入)
// 但为了确保流程可以继续,我们仍然标记为成功
state.setProcessStatus(NODE_PROCESS_STATUS_SUCCESS);
log.info("No human feedback found for node: {}, continuing workflow", node.getTitle());
}
return new NodeProcessResult();
}
}

View File

@@ -1,106 +0,0 @@
package org.ruoyi.workflow.workflow.node.keywordExtractor;
import dev.langchain4j.data.message.SystemMessage;
import lombok.extern.slf4j.Slf4j;
import org.apache.commons.lang3.StringUtils;
import org.ruoyi.workflow.entity.WorkflowComponent;
import org.ruoyi.workflow.entity.WorkflowNode;
import org.ruoyi.workflow.util.SpringUtil;
import org.ruoyi.workflow.util.WorkflowMessageUtil;
import org.ruoyi.workflow.workflow.NodeProcessResult;
import org.ruoyi.workflow.workflow.WfNodeState;
import org.ruoyi.workflow.workflow.WfState;
import org.ruoyi.workflow.workflow.WorkflowUtil;
import org.ruoyi.workflow.workflow.data.NodeIOData;
import org.ruoyi.workflow.workflow.node.AbstractWfNode;
import org.ruoyi.workflow.workflow.node.enmus.NodeMessageTemplateEnum;
import java.util.ArrayList;
import java.util.List;
import static org.ruoyi.workflow.cosntant.AdiConstant.WorkflowConstant.DEFAULT_OUTPUT_PARAM_NAME;
/**
* 【节点】关键词提取节点
* 使用 LLM 从文本中提取关键词
*/
@Slf4j
public class KeywordExtractorNode extends AbstractWfNode {
public KeywordExtractorNode(WorkflowComponent wfComponent, WorkflowNode nodeDef, WfState wfState, WfNodeState nodeState) {
super(wfComponent, nodeDef, wfState, nodeState);
}
/**
* 处理关键词提取
* nodeConfig 格式:
* {
* "model_name": "deepseek-chat",
* "category": "llm",
* "top_n": 5,
* "prompt": "额外的提示词"
* }
*
* @return 提取的关键词列表
*/
@Override
public NodeProcessResult onProcess() {
KeywordExtractorNodeConfig config = checkAndGetConfig(KeywordExtractorNodeConfig.class);
// 获取输入文本
String inputText = getFirstInputText();
if (StringUtils.isBlank(inputText)) {
log.warn("Keyword extractor node has no input text, node: {}", state.getUuid());
// 返回空结果
List<NodeIOData> outputs = new ArrayList<>();
outputs.add(NodeIOData.createByText(DEFAULT_OUTPUT_PARAM_NAME, "", ""));
return NodeProcessResult.builder().content(outputs).build();
}
log.info("Keyword extractor node config: {}", config);
log.info("Input text length: {}", inputText.length());
// 构建提示词
String prompt = buildPrompt(config, inputText);
log.info("Keyword extraction prompt: {}", prompt);
// 调用 LLM 进行关键词提取
WorkflowUtil workflowUtil = SpringUtil.getBean(WorkflowUtil.class);
String modelName = config.getModelName();
// 获取节点模板提示词信息
String nodeMessageTemplate = WorkflowMessageUtil.getNodeMessageTemplate(NodeMessageTemplateEnum.KEYWORD_EXTRACTOR.getValue());
// 发送SSE事件消息
WorkflowMessageUtil.sendEmitterMessage(wfState.getSseEmitter(), node, nodeMessageTemplate);
// 使用流式调用
workflowUtil.streamingInvokeLLM(wfState, state, node, modelName, prompt, nodeMessageTemplate);
return new NodeProcessResult();
}
/**
* 构建关键词提取的提示词
*/
private String buildPrompt(KeywordExtractorNodeConfig config, String inputText) {
StringBuilder promptBuilder = new StringBuilder();
// 基础提示词
promptBuilder.append("请从以下文本中提取 ").append(config.getTopN()).append(" 个最重要的关键词。\n\n");
// 添加自定义提示词(如果有)
if (StringUtils.isNotBlank(config.getPrompt())) {
promptBuilder.append(config.getPrompt()).append("\n\n");
}
// 输出格式要求
promptBuilder.append("要求:\n");
promptBuilder.append("1. 只返回关键词,每个关键词用逗号分隔\n");
promptBuilder.append("2. 关键词应该是名词或名词短语\n");
promptBuilder.append("3. 按重要性从高到低排序\n");
promptBuilder.append("4. 不要添加任何解释或额外的文字\n\n");
// 原始文本
promptBuilder.append("文本内容:\n");
promptBuilder.append(inputText);
return promptBuilder.toString();
}
}

View File

@@ -1,42 +0,0 @@
package org.ruoyi.workflow.workflow.node.keywordExtractor;
import com.fasterxml.jackson.annotation.JsonProperty;
import jakarta.validation.constraints.Max;
import jakarta.validation.constraints.Min;
import jakarta.validation.constraints.NotNull;
import lombok.Data;
import lombok.EqualsAndHashCode;
/**
* 关键词提取节点配置
*/
@EqualsAndHashCode
@Data
public class KeywordExtractorNodeConfig {
/**
* 模型分类llm, embedding 等)
*/
private String category;
/**
* 模型名称
*/
@NotNull
@JsonProperty("model_name")
private String modelName;
/**
* 提取的关键词数量
*/
@Min(1)
@Max(50)
@JsonProperty("top_n")
private Integer topN = 5;
/**
* 提示词(可选)
* 用于指导关键词提取的额外说明
*/
private String prompt;
}

View File

@@ -7,6 +7,7 @@ import lombok.extern.slf4j.Slf4j;
import org.apache.commons.lang3.StringUtils;
import org.ruoyi.workflow.entity.WorkflowComponent;
import org.ruoyi.workflow.entity.WorkflowNode;
import org.ruoyi.workflow.util.JsonUtil;
import org.ruoyi.workflow.workflow.NodeProcessResult;
import org.ruoyi.workflow.workflow.WfNodeState;
import org.ruoyi.workflow.workflow.WfState;
@@ -40,15 +41,25 @@ public class MailSendNode extends AbstractWfNode {
String input = getDataFromInput(inputs);
// 判断是否为JSON格式(LLM输出转换 由LLM生成格式)
if (StringUtils.isNotBlank(input) && isJson(input)) {
// 使用Jackson解析和合并配置
ObjectMapper objectMapper = new ObjectMapper();
JsonNode inputJson = objectMapper.readTree(input);
// 将config转换为JsonNode
JsonNode configJson = objectMapper.valueToTree(config);
// 合并两个JSON节点
JsonNode mergedJson = objectMapper.readerForUpdating(configJson).readValue(inputJson);
// 转换回config对象
config = objectMapper.treeToValue(mergedJson, MailSendNodeConfig.class);
try {
// 使用统一的 JsonUtil 进行解析
JsonNode inputJson = JsonUtil.toJsonNode(input);
if (inputJson != null) {
// 使用 JsonUtil 内部的 ObjectMapper 进行合并
ObjectMapper objectMapper = new ObjectMapper();
// 将config转换为JsonNode
JsonNode configJson = objectMapper.valueToTree(config);
// 合并两个JSON节点
JsonNode mergedJson = objectMapper.readerForUpdating(configJson).readValue(inputJson);
// 转换回config对象
config = objectMapper.treeToValue(mergedJson, MailSendNodeConfig.class);
} else {
log.warn("输入 JSON 解析结果为 null使用原始配置");
}
} catch (Exception e) {
log.error("合并邮件配置失败,使用原始配置: {}", e.getMessage(), e);
// 继续使用原始 config不中断流程
}
}
// 安全获取模板(使用 defaultString 避免 null

View File

@@ -1,7 +1,6 @@
package org.ruoyi.workflow.workflow.node.switcher;
import com.fasterxml.jackson.databind.JsonNode;
import com.fasterxml.jackson.databind.ObjectMapper;
import lombok.extern.slf4j.Slf4j;
import org.apache.commons.lang3.ObjectUtils;
import org.apache.commons.lang3.StringUtils;
@@ -9,6 +8,7 @@ import org.ruoyi.common.core.utils.SpringUtils;
import org.ruoyi.workflow.entity.WorkflowComponent;
import org.ruoyi.workflow.entity.WorkflowNode;
import org.ruoyi.workflow.service.WorkflowNodeService;
import org.ruoyi.workflow.util.JsonUtil;
import org.ruoyi.workflow.workflow.NodeProcessResult;
import org.ruoyi.workflow.workflow.WfNodeState;
import org.ruoyi.workflow.workflow.WfState;
@@ -18,8 +18,6 @@ import org.ruoyi.workflow.workflow.node.enmus.NodeMessageTemplateEnum;
import java.math.BigDecimal;
import java.util.List;
import java.util.Objects;
import java.util.Optional;
/**
* 条件分支节点
@@ -339,9 +337,12 @@ public class SwitcherNode extends AbstractWfNode {
log.info("节点 '{}' 的输入配置: {}", nodeUuid, inputConfig);
if (StringUtils.isNotBlank(inputConfig)){
try {
// 使用Jackson解析输入配置
ObjectMapper objectMapper = new ObjectMapper();
JsonNode configJson = objectMapper.readTree(inputConfig);
// 使用统一的 JsonUtil 而不是每次创建新的 ObjectMapper
JsonNode configJson = JsonUtil.toJsonNode(inputConfig);
if (configJson == null) {
log.warn("节点 '{}' 的输入配置 JSON 解析结果为 null", nodeUuid);
return result;
}
// 获取 user_inputs 数组
JsonNode userInputs = configJson.get("user_inputs");
if (userInputs != null && userInputs.isArray()) {
@@ -357,7 +358,8 @@ public class SwitcherNode extends AbstractWfNode {
}
}
} catch (Exception e) {
log.error("解析节点输入配置失败: {}", nodeUuid, e);
log.error("解析节点 '{}' 输入配置失败,参数名: {}, 配置内容: {}", nodeUuid, paramName, inputConfig, e);
// 不抛出异常,返回默认结果,避免中断整个流程
}
}
}

View File

@@ -0,0 +1,21 @@
-- =============================================
-- 流程编排搜索节点配置脚本
-- =============================================
-- 说明:本脚本用于添加搜索节点的系统配置
-- 执行前请确保 sys_config 表存在
-- =============================================
-- 搜索节点模板配置
INSERT INTO sys_config (config_name, config_key, config_value, config_type, remark, create_by, create_time, update_by, update_time)
VALUES ('搜索节点模板', 'node.googleSearch.template', '正在搜索相关内容...', 'Y', '搜索节点的响应模板,用于网络搜索功能', 'admin', NOW(), 'admin', NOW())
ON DUPLICATE KEY UPDATE config_value = '正在搜索相关内容...', update_time = NOW();
-- =============================================
-- 验证配置是否添加成功
-- =============================================
SELECT config_id, config_name, config_key, config_value, config_type, remark
FROM sys_config
WHERE config_key IN (
'node.googleSearch.template'
)
ORDER BY config_id;

View File

@@ -1,425 +0,0 @@
# Ruoyi-AI 流程编排模块详细说明文档
## 概述
Ruoyi-AI 工作流模块是一个基于 LangGraph4j 的智能工作流引擎支持可视化工作流设计、AI 模型集成、条件分支、人机交互等高级功能。该模块采用微服务架构,提供完整的
RESTful API 和流式响应支持。
## 模块架构
### 1. 核心依赖
- **LangGraph4j**: 1.5.3 - 工作流图执行引擎
- **LangChain4j**: 1.11.0 - AI 模型集成框架
- **Spring Boot**: 3.5.8 - 应用框架
- **MyBatis Plus**: 数据访问层
- **Redis**: 缓存和状态管理
- **OpenAPI**: API 文档
## 核心功能
### 1. 工作流管理
#### 1.1 工作流定义
- **创建工作流**: 支持自定义标题、描述、公开性设置
- **编辑工作流**: 可视化节点编辑、连接线配置
- **版本控制**: 支持工作流的版本管理和回滚
- **权限管理**: 支持公开/私有工作流设置
#### 1.2 工作流执行
- **流式执行**: 基于 SSE 的实时流式响应
- **状态管理**: 完整的执行状态跟踪
- **错误处理**: 详细的错误信息和异常处理
- **中断恢复**: 支持工作流中断和恢复执行
### 2. 节点类型
#### 2.1 基础节点
- **Start**: 开始节点,定义工作流入口
- **End**: 结束节点,定义工作流出口
#### 2.2 AI 模型节点
- **Answer**: 大语言模型问答节点
- **Dalle3**: DALL-E 3 图像生成
- **Tongyiwanx**: 通义万相图像生成
- **Classifier**: 内容分类节点
#### 2.3 数据处理节点
- **DocumentExtractor**: 文档信息提取
- **KeywordExtractor**: 关键词提取
- **FaqExtractor**: 常见问题提取
- **KnowledgeRetrieval**: 知识库检索
#### 2.4 控制流节点
- **Switcher**: 条件分支节点
- **HumanFeedback**: 人机交互节点
#### 2.5 外部集成节点
- **Google**: Google 搜索集成
- **MailSend**: 邮件发送
- **HttpRequest**: HTTP 请求
- **Template**: 模板转换
### 3. 数据流管理
#### 3.1 输入输出定义
```java
// 节点输入输出数据结构
public class NodeIOData {
private String name; // 参数名称
private NodeIODataContent content; // 参数内容
}
// 支持的数据类型
public enum WfIODataTypeEnum {
TEXT, // 文本
NUMBER, // 数字
BOOLEAN, // 布尔值
FILES, // 文件
OPTIONS // 选项
}
```
#### 3.2 参数引用
- **节点间引用**: 支持上游节点输出作为下游节点输入
- **参数映射**: 自动处理参数名称映射
- **类型转换**: 自动进行数据类型转换
## 数据库设计
### 1. 核心表结构
#### 1.1 工作流定义表 (t_workflow)
```sql
CREATE TABLE t_workflow (
id BIGINT AUTO_INCREMENT PRIMARY KEY,
uuid VARCHAR(32) NOT NULL DEFAULT '',
title VARCHAR(100) NOT NULL DEFAULT '',
remark TEXT NOT NULL DEFAULT '',
user_id BIGINT NOT NULL DEFAULT 0,
is_public TINYINT(1) NOT NULL DEFAULT 0,
is_enable TINYINT(1) NOT NULL DEFAULT 1,
create_time DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP,
update_time DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP,
is_deleted TINYINT(1) NOT NULL DEFAULT 0
);
```
#### 1.2 工作流节点表 (t_workflow_node)
```sql
CREATE TABLE t_workflow_node (
id BIGINT AUTO_INCREMENT PRIMARY KEY,
uuid VARCHAR(32) NOT NULL DEFAULT '',
workflow_id BIGINT NOT NULL DEFAULT 0,
workflow_component_id BIGINT NOT NULL DEFAULT 0,
user_id BIGINT NOT NULL DEFAULT 0,
title VARCHAR(100) NOT NULL DEFAULT '',
remark VARCHAR(500) NOT NULL DEFAULT '',
input_config JSON NOT NULL DEFAULT ('{}'),
node_config JSON NOT NULL DEFAULT ('{}'),
position_x DOUBLE NOT NULL DEFAULT 0,
position_y DOUBLE NOT NULL DEFAULT 0,
create_time DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP,
update_time DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP,
is_deleted TINYINT(1) NOT NULL DEFAULT 0
);
```
#### 1.3 工作流边表 (t_workflow_edge)
```sql
CREATE TABLE t_workflow_edge (
id BIGINT AUTO_INCREMENT PRIMARY KEY,
uuid VARCHAR(32) NOT NULL DEFAULT '',
workflow_id BIGINT NOT NULL DEFAULT 0,
source_node_uuid VARCHAR(32) NOT NULL DEFAULT '',
source_handle VARCHAR(32) NOT NULL DEFAULT '',
target_node_uuid VARCHAR(32) NOT NULL DEFAULT '',
create_time DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP,
update_time DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP,
is_deleted TINYINT(1) NOT NULL DEFAULT 0
);
```
#### 1.4 工作流运行时表 (t_workflow_runtime)
```sql
CREATE TABLE t_workflow_runtime (
id BIGINT AUTO_INCREMENT PRIMARY KEY,
uuid VARCHAR(32) NOT NULL DEFAULT '',
user_id BIGINT NOT NULL DEFAULT 0,
workflow_id BIGINT NOT NULL DEFAULT 0,
input JSON NOT NULL DEFAULT ('{}'),
output JSON NOT NULL DEFAULT ('{}'),
status SMALLINT NOT NULL DEFAULT 1,
status_remark VARCHAR(250) NOT NULL DEFAULT '',
create_time DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP,
update_time DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP,
is_deleted TINYINT(1) NOT NULL DEFAULT 0
);
```
#### 1.5 工作流组件表 (t_workflow_component)
```sql
CREATE TABLE t_workflow_component (
id BIGINT AUTO_INCREMENT PRIMARY KEY,
uuid VARCHAR(32) DEFAULT '' NOT NULL,
name VARCHAR(32) DEFAULT '' NOT NULL,
title VARCHAR(100) DEFAULT '' NOT NULL,
remark TEXT NOT NULL,
display_order INT DEFAULT 0 NOT NULL,
is_enable TINYINT(1) DEFAULT 0 NOT NULL,
create_time DATETIME DEFAULT CURRENT_TIMESTAMP NOT NULL,
update_time DATETIME DEFAULT CURRENT_TIMESTAMP NOT NULL,
is_deleted TINYINT(1) DEFAULT 0 NOT NULL
);
```
## API 接口
### 1. 工作流管理接口
#### 1.1 基础操作
```http
#
POST /workflow/add
Content-Type: application/json
{
"title": "",
"remark": "",
"isPublic": false
}
#
POST /workflow/update
Content-Type: application/json
{
"uuid": "UUID",
"title": "",
"remark": ""
}
#
POST /workflow/del/{uuid}
# /
POST /workflow/enable/{uuid}?enable=true
```
#### 1.2 搜索和查询
```http
#
GET /workflow/mine/search?keyword=&isPublic=true&currentPage=1&pageSize=10
#
GET /workflow/public/search?keyword=&currentPage=1&pageSize=10
#
GET /workflow/public/component/list
```
### 2. 工作流执行接口
#### 2.1 流式执行
```http
#
POST /workflow/run
Content-Type: application/json
Accept: text/event-stream
{
"uuid": "UUID",
"inputs": [
{
"name": "input",
"content": {
"type": 1,
"textContent": ""
}
}
]
}
```
#### 2.2 运行时管理
```http
#
POST /workflow/runtime/resume/{runtimeUuid}
Content-Type: application/json
{
"feedbackContent": ""
}
#
GET /workflow/runtime/page?wfUuid=UUID&currentPage=1&pageSize=10
#
GET /workflow/runtime/nodes/{runtimeUuid}
#
POST /workflow/runtime/clear?wfUuid=UUID
```
### 3. 管理端接口
#### 3.1 工作流管理
```http
#
POST /admin/workflow/search
Content-Type: application/json
{
"title": "",
"isPublic": true,
"isEnable": true
}
# /
POST /admin/workflow/enable?uuid=UUID&isEnable=true
```
## 核心实现
### 1. 工作流引擎 (WorkflowEngine)
工作流引擎是整个模块的核心,负责:
- 工作流图的构建和编译
- 节点执行调度
- 状态管理和持久化
- 流式输出处理
```java
public class WorkflowEngine {
// 核心执行方法
public void run(User user, List<ObjectNode> userInputs, SseEmitter sseEmitter) {
// 1. 验证工作流状态
// 2. 创建运行时实例
// 3. 构建状态图
// 4. 执行工作流
// 5. 处理流式输出
}
// 恢复执行方法
public void resume(String userInput) {
// 1. 更新状态
// 2. 继续执行
}
}
```
### 2. 节点工厂 (WfNodeFactory)
节点工厂负责根据组件类型创建对应的节点实例:
```java
public class WfNodeFactory {
public static AbstractWfNode create(WorkflowComponent component,
WorkflowNode node,
WfState wfState,
WfNodeState nodeState) {
// 根据组件类型创建对应的节点实例
switch (component.getName()) {
case "Answer":
return new LLMAnswerNode(component, node, wfState, nodeState);
case "Switcher":
return new SwitcherNode(component, node, wfState, nodeState);
// ... 其他节点类型
}
}
}
```
### 3. 图构建器 (WorkflowGraphBuilder)
图构建器负责将工作流定义转换为可执行的状态图:
```java
public class WorkflowGraphBuilder {
public StateGraph<WfNodeState> build(WorkflowNode startNode) {
// 1. 构建编译节点树
// 2. 转换为状态图
// 3. 添加节点和边
// 4. 处理条件分支
// 5. 处理并行执行
}
}
```
## 流式响应机制
### 1. SSE 事件类型
工作流执行过程中会发送多种类型的 SSE 事件:
```javascript
// 节点开始执行
[NODE_RUN_节点UUID] - 节点执行开始事件
// 节点输入数据
[NODE_INPUT_节点UUID] - 节点输入数据事件
// 节点输出数据
[NODE_OUTPUT_节点UUID] - 节点输出数据事件
// 流式内容块
[NODE_CHUNK_节点UUID] - 流式内容块事件
// 等待用户输入
[NODE_WAIT_FEEDBACK_BY_节点UUID] - 等待用户输入事件
```
### 2. 流式处理流程
1. **初始化**: 创建工作流运行时实例
2. **节点执行**: 逐个执行工作流节点
3. **实时输出**: 通过 SSE 实时推送执行结果
4. **状态更新**: 实时更新节点和工作流状态
5. **错误处理**: 捕获并处理执行过程中的错误
## 扩展开发
### 1. 自定义节点开发
要开发自定义工作流节点,需要:
1. **创建节点类**:继承 `AbstractWfNode`
2. **实现处理逻辑**:重写 `onProcess()` 方法
3. **定义配置类**:创建节点配置类
4. **注册组件**:在组件表中注册新组件
```java
public class CustomNode extends AbstractWfNode {
@Override
protected NodeProcessResult onProcess() {
// 实现自定义处理逻辑
List<NodeIOData> outputs = new ArrayList<>();
// ... 处理逻辑
return NodeProcessResult.success(outputs);
}
}
```
### 2. 自定义组件注册
```sql
-- 在 t_workflow_component 表中添加新组件
INSERT INTO t_workflow_component (uuid, name, title, remark, is_enable)
VALUES (REPLACE(UUID(), '-', ''), 'CustomNode', '自定义节点', '自定义节点描述', true);
```

View File

@@ -1,10 +1,8 @@
package org.ruoyi.controller.chat;
import jakarta.servlet.http.HttpServletRequest;
import jakarta.validation.Valid;
import lombok.RequiredArgsConstructor;
import lombok.extern.slf4j.Slf4j;
import org.ruoyi.common.chat.domain.dto.request.AgentChatRequest;
import org.ruoyi.common.chat.domain.dto.request.ChatRequest;
import org.ruoyi.service.chat.impl.ChatServiceFacade;
import org.springframework.stereotype.Controller;

View File

@@ -150,35 +150,59 @@ public class ChatServiceFacade implements IChatService {
* @return SseEmitter
*/
public SseEmitter sseChat(ChatRequest chatRequest) {
// 具体的服务实现
Long userId = LoginHelper.getUserId();
String tokenValue = StpUtil.getTokenValue();
// 每个会话一个 SSE 连接,避免同用户多会话串台
SseEmitter emitter = sseEmitterManager.connect(String.valueOf(chatRequest.getSessionId()));
boolean workflowMode = Boolean.TRUE.equals(chatRequest.getEnableWorkFlow());
boolean agentMode = chatRequest.getAgentId() != null;
if (workflowMode && agentMode) {
throw new IllegalArgumentException("对话模式参数冲突:工作流和智能体不能同时启用");
}
// 工作流模式。工作流引擎负责创建并持有自己的 SSE必须在普通聊天连接创建前路由。
if (workflowMode) {
chatMessageService.saveChatMessage(
userId,
chatRequest.getSessionId(),
chatRequest.getContent(),
RoleType.USER.getName(),
chatRequest.getModel()
);
return handleWorkflowChat(chatRequest);
}
// 智能体解析:传入 agentId 时按智能体绑定的模型覆盖 model 字段
AgentVo agentVo = null;
if (chatRequest.getAgentId() != null) {
if (agentMode) {
agentVo = agentService.queryById(chatRequest.getAgentId());
if (agentVo == null) {
throw new IllegalArgumentException("智能体不存在: " + chatRequest.getAgentId());
}
if (agentVo != null && agentVo.getModelId() != null) {
ChatModelVo agentModel = chatModelService.queryById(agentVo.getModelId());
if (agentModel != null) {
chatRequest.setModel(agentModel.getModelName());
if (agentModel == null) {
throw new IllegalArgumentException("智能体绑定的模型不存在: " + agentVo.getModelId());
}
} else {
log.warn("智能体不存在或未配置模型,回退到 model 字段: agentId={}", chatRequest.getAgentId());
chatRequest.setModel(agentModel.getModelName());
}
}
if (StringUtils.isBlank(chatRequest.getModel())) {
throw new IllegalArgumentException(
agentVo == null ? "对话模式必须指定模型" : "智能体未绑定模型,且请求未提供回退模型"
);
}
// 根据模型名称查询完整配置
ChatModelVo chatModelVo = chatModelService.selectModelByName(chatRequest.getModel());
if (chatModelVo == null) {
throw new IllegalArgumentException("模型不存在: " + chatRequest.getModel());
}
// 对话和智能体模式共用按会话隔离的 SSE。
SseEmitter emitter = sseEmitterManager.connect(String.valueOf(chatRequest.getSessionId()));
// 构建上下文消息列表(系统提示词 + 历史消息 + 当前用户消息)
// 注意RAG 检索增强统一在 handleAgentChat 中执行一次,此处不再重复检索
List<ChatMessage> contextMessages = buildContextMessages(chatRequest, agentVo);
chatRequest.setEmitter(emitter);
@@ -190,43 +214,68 @@ public class ChatServiceFacade implements IChatService {
// 保存用户消息
chatMessageService.saveChatMessage(userId, chatRequest.getSessionId(), chatRequest.getContent(), RoleType.USER.getName(), chatRequest.getModel());
TraceRunHandle traceRun = Boolean.TRUE.equals(chatRequest.getEnableWorkFlow())
? null : startRagTraceRun(chatRequest, userId);
// 3. 路由对话模式:工作流对话 / 智能体对话(两者均返回各自的 SseEmitter
return handleSpecialChatModes(chatRequest, agentVo, traceRun);
}
TraceRunHandle traceRun = startRagTraceRun(chatRequest, userId);
/**
* 路由对话模式:仅两种情况——工作流对话 / 智能体对话。
*
* @param chatRequest 聊天请求
* @param agentVo 智能体配置(可为 null
* @return 对应模式的 SseEmitter
*/
private SseEmitter handleSpecialChatModes(ChatRequest chatRequest, AgentVo agentVo, TraceRunHandle traceRun) {
// 模式1工作流对话前端应用市场选工作流后携带 workFlowRunner
if (Boolean.TRUE.equals(chatRequest.getEnableWorkFlow())) {
log.info("处理工作流对话,会话: {}", chatRequest.getSessionId());
WorkFlowRunner runner = chatRequest.getWorkFlowRunner();
if (ObjectUtils.isEmpty(runner)) {
log.warn("工作流参数为空");
}
return workFlowStarterService.streaming(
ThreadContext.getCurrentUser(),
runner.getUuid(),
runner.getInputs(),
chatRequest.getSessionId()
);
// 智能体和普通对话互斥:有 agentId 为智能体,否则为普通模型对话。
if (agentVo != null) {
log.info("处理智能体对话,会话:{},agentId:{}", chatRequest.getSessionId(), chatRequest.getAgentId());
return handleAgentChat(chatRequest, agentVo, traceRun);
}
// 模式2智能体对话默认走 Supervisor 多 Agent 编排)
return handleAgentChat(chatRequest, agentVo, traceRun);
log.info("处理普通对话,会话:{},模型:{}", chatRequest.getSessionId(), chatRequest.getModel());
return handleModelChat(chatRequest, traceRun);
}
/**
* 智能体对话模式(默认):构建 Supervisor 多 Agent 编排并异步执行,结果通过 SSE 推送
* 工作流模式。工作流运行时负责 SSE、节点执行和结束事件
*/
private SseEmitter handleWorkflowChat(ChatRequest chatRequest) {
WorkFlowRunner runner = chatRequest.getWorkFlowRunner();
if (ObjectUtils.isEmpty(runner) || StringUtils.isBlank(runner.getUuid())) {
throw new IllegalArgumentException("工作流模式必须提供 workFlowRunner.uuid");
}
log.info("处理工作流对话,会话:{},workflowUuid:{}", chatRequest.getSessionId(), runner.getUuid());
return workFlowStarterService.streaming(
ThreadContext.getCurrentUser(),
runner.getUuid(),
runner.getInputs() == null ? List.of() : runner.getInputs(),
chatRequest.getSessionId()
);
}
/**
* 普通对话模式:直接调用选定模型,不装配 Supervisor、MCP、Skills 或专业子 Agent。
*/
private SseEmitter handleModelChat(ChatRequest chatRequest, TraceRunHandle traceRun) {
ChatModelVo chatModelVo = chatRequest.getChatModelVo();
AbstractChatService chatService = chatServiceFactory.getOriginalService(chatModelVo.getProviderCode());
StreamingChatModel streamingChatModel = chatService.buildStreamingChatModel(chatModelVo, chatRequest);
List<ChatMessage> messages = buildModelChatMessages(chatRequest);
TraceStreamSpan llmSpan = null;
try (TraceScope ignored = openTraceScope(traceRun, chatRequest.getUserId())) {
llmSpan = startLlmCallSpan(traceRun, chatRequest, "handleModelChat");
streamingChatModel.chat(
messages,
createModelChatResponseHandler(chatRequest, traceRun, llmSpan)
);
} catch (Exception e) {
if (llmSpan != null) {
llmSpan.finishError(e);
llmSpan.detach();
}
finishTraceRun(traceRun, TraceConstants.STATUS_ERROR, e);
SseMessageUtils.sendError(String.valueOf(chatRequest.getSessionId()), e.getMessage());
SseMessageUtils.completeConnection(String.valueOf(chatRequest.getSessionId()));
log.error("普通对话执行失败", e);
}
return chatRequest.getEmitter();
}
/**
* 智能体对话模式:构建 Supervisor 多 Agent 编排并异步执行,结果通过 SSE 推送。
*
* @param chatRequest 聊天请求
* @param agentVo 智能体配置(可为 null无智能体时用请求 model 兜底)
* @param agentVo 智能体配置
*/
private SseEmitter handleAgentChat(ChatRequest chatRequest, AgentVo agentVo, TraceRunHandle traceRun) {
ChatModelVo chatModelVo = chatRequest.getChatModelVo();
@@ -321,7 +370,7 @@ public class ChatServiceFacade implements IChatService {
CompletableFuture.runAsync(() -> {
TraceStreamSpan llmSpan = null;
try (TraceScope ignored = openTraceScope(traceRun, userId)) {
llmSpan = startLlmCallSpan(traceRun, chatRequest);
llmSpan = startLlmCallSpan(traceRun, chatRequest, "handleAgentChat");
String result = supervisor.invoke(prompt);
SseMessageUtils.sendContent(sessionId, result);
SseMessageUtils.sendDone(sessionId);
@@ -386,7 +435,8 @@ public class ChatServiceFacade implements IChatService {
traceRun.businessId, userId, traceRun.tenantId);
}
private TraceStreamSpan startLlmCallSpan(TraceRunHandle traceRun, ChatRequest chatRequest) {
private TraceStreamSpan startLlmCallSpan(TraceRunHandle traceRun, ChatRequest chatRequest,
String methodName) {
if (traceRun == null || StringUtils.isBlank(TraceContext.getTraceId())) {
return null;
}
@@ -401,7 +451,7 @@ public class ChatServiceFacade implements IChatService {
node.setNodeName("llm-call");
node.setNodeType(RagTraceNodeTypes.NODE_LLM_CALL);
node.setClassName(ChatServiceFacade.class.getName());
node.setMethodName("handleAgentChat");
node.setMethodName(methodName);
node.setStatus(TraceConstants.STATUS_RUNNING);
node.setStartTime(new Date(startMillis));
node.setInputPayload(RagTracePayloadBuilder.streamInputSummary(chatRequest));
@@ -564,7 +614,11 @@ public class ChatServiceFacade implements IChatService {
Long userId = LoginHelper.getUserId();
// 5. 建立 SSE 连接(用于前端监听,按会话隔离)
sseEmitterManager.connect(String.valueOf(chatRequest.getSessionId()));
// 工作流调用时(externalHandler 非空), SSE 连接由工作流引擎创建并持有(WorkflowStarter#streaming),
// connect 为替换语义(关闭同键旧连接), 此处重连会掐断工作流连接, 必须跳过
if (externalHandler == null) {
sseEmitterManager.connect(String.valueOf(chatRequest.getSessionId()));
}
// 保存用户消息
chatMessageService.saveChatMessage(userId, chatRequest.getSessionId(), chatRequest.getContent(), RoleType.USER.getName(), chatRequest.getModel());
@@ -645,6 +699,19 @@ public class ChatServiceFacade implements IChatService {
return messages;
}
/**
* 构建普通对话消息。保留历史上下文,并在请求指定知识库时仅增强当前用户消息。
*/
private List<ChatMessage> buildModelChatMessages(ChatRequest chatRequest) {
List<ChatMessage> messages = new ArrayList<>(chatRequest.getContextMessages());
String augmentedInput = augmentAgentInput(chatRequest, null);
int lastIndex = messages.size() - 1;
if (lastIndex >= 0 && messages.get(lastIndex) instanceof UserMessage) {
messages.set(lastIndex, UserMessage.userMessage(augmentedInput));
}
return messages;
}
/**
* 将上下文消息格式化为多轮对话文本(供只接受 String 输入的 Supervisor 使用)。
* 跳过 SystemMessage系统提示词单独前置与最后一条当前用户消息单独做 RAG 增强后拼接)。
@@ -791,6 +858,77 @@ public class ChatServiceFacade implements IChatService {
return queryVectorBo;
}
/**
* 普通对话响应处理器:推送流式内容、保存助手消息并结束链路追踪。
*/
private StreamingChatResponseHandler createModelChatResponseHandler(ChatRequest chatRequest,
TraceRunHandle traceRun,
TraceStreamSpan llmSpan) {
String sessionId = String.valueOf(chatRequest.getSessionId());
return new StreamingChatResponseHandler() {
private final StringBuilder messageBuffer = new StringBuilder();
@Override
public void onPartialResponse(String partialResponse) {
messageBuffer.append(partialResponse);
SseMessageUtils.sendContent(sessionId, partialResponse);
}
@Override
public void onPartialThinking(PartialThinking partialThinking) {
SseMessageUtils.sendReasoning(sessionId, partialThinking.text());
}
@Override
public void onCompleteResponse(ChatResponse completeResponse) {
try {
String fullMessage = messageBuffer.toString();
if (StringUtils.isNotBlank(fullMessage)) {
chatMessageService.saveChatMessage(
chatRequest.getUserId(),
chatRequest.getSessionId(),
fullMessage,
RoleType.ASSISTANT.getName(),
chatRequest.getModel()
);
} else {
log.warn("普通对话返回空消息,会话:{}", chatRequest.getSessionId());
}
if (llmSpan != null) {
llmSpan.finishSuccess(RagTracePayloadBuilder.streamOutputSummary(fullMessage.length()));
}
finishTraceRun(traceRun, TraceConstants.STATUS_SUCCESS, null);
SseMessageUtils.sendDone(sessionId);
} catch (Exception e) {
if (llmSpan != null) {
llmSpan.finishError(e);
}
finishTraceRun(traceRun, TraceConstants.STATUS_ERROR, e);
SseMessageUtils.sendError(sessionId, e.getMessage());
log.error("普通对话完成处理失败", e);
} finally {
if (llmSpan != null) {
llmSpan.detach();
}
SseMessageUtils.completeConnection(sessionId);
}
}
@Override
public void onError(Throwable error) {
if (llmSpan != null) {
llmSpan.finishError(error);
llmSpan.detach();
}
finishTraceRun(traceRun, TraceConstants.STATUS_ERROR, error);
SseMessageUtils.sendError(sessionId, error.getMessage());
SseMessageUtils.completeConnection(sessionId);
log.error("普通对话流式响应失败", error);
}
};
}
/**
* 创建组合响应处理器 - 同时发送到 SSE 和外部 handler
*
@@ -811,7 +949,10 @@ public class ChatServiceFacade implements IChatService {
messageBuffer.append(partialResponse);
// 2. 发送内容事件到 SSE前端可通过 SSE 监听)
SseMessageUtils.sendContent(sessionId, partialResponse);
// 工作流调用时连接归工作流引擎所有, token 由引擎以 [NODE_CHUNK_] 事件推送, 不走聊天协议
if (externalHandler == null) {
SseMessageUtils.sendContent(sessionId, partialResponse);
}
// 3. 转发给外部 handlerWorkflow 等模块可处理)
if (externalHandler != null) {
@@ -821,8 +962,10 @@ public class ChatServiceFacade implements IChatService {
@Override
public void onPartialThinking(PartialThinking partialThinking) {
// 发送推理内容到 SSE前端通过 reasoning 事件监听)
SseMessageUtils.sendReasoning(sessionId, partialThinking.text());
// 发送推理内容到 SSE前端通过 reasoning 事件监听), 工作流调用时不发送
if (externalHandler == null) {
SseMessageUtils.sendReasoning(sessionId, partialThinking.text());
}
// 转发给外部 handler
if (externalHandler != null) {
@@ -833,11 +976,12 @@ public class ChatServiceFacade implements IChatService {
@Override
public void onCompleteResponse(ChatResponse completeResponse) {
try {
// 1. 发送完成事件
SseMessageUtils.sendDone(sessionId);
// 2. 关闭 SSE 连接
SseMessageUtils.completeConnection(sessionId);
// 1&2. 发送完成事件并关闭 SSE 连接
// 工作流调用时流程可能还有后续节点, 连接关闭由工作流引擎统一负责, 此处不能关闭
if (externalHandler == null) {
SseMessageUtils.sendDone(sessionId);
SseMessageUtils.completeConnection(sessionId);
}
// 3. 转发给外部 handler
if (externalHandler != null) {
@@ -850,8 +994,10 @@ public class ChatServiceFacade implements IChatService {
@Override
public void onError(Throwable error) {
// 发送错误事件
SseMessageUtils.sendError(sessionId, error.getMessage());
// 发送错误事件(工作流调用时由工作流引擎统一上报)
if (externalHandler == null) {
SseMessageUtils.sendError(sessionId, error.getMessage());
}
log.error("流式响应错误: {}", error.getMessage(), error);
// 转发给外部 handler