Compare commits

...

13 Commits

Author SHA1 Message Date
4e5c17d59f Merge remote-tracking branch 'origin/dev' into dev 2026-05-16 13:55:27 +08:00
79da694baf feat(flow): 添加任务执行入口和实例管理功能
- 在 FlowTaskRuntimeEntry 中添加事务注解支持
- 为 FlowTaskRuntimeService 接口添加方法间隔
- 在 TeFlowController 中引入任务实例ID列表管理
- 新增 demo 接口用于任务执行和实例清理
- 实现任务实例的自动停止和清理机制
- 集成 flowControlService 的 stop 功能用于流程控制
2026-05-16 13:55:20 +08:00
3000531ca9 ```
fix(config):修复生产环境配置中的拼写错误修正了 minio服务地址配置中 "mini" 到 "minio" 的拼写错误,
以及 API 密钥末尾多余的空白字符。

这两个配置项都位于 application-prod.yml 文件中,确保了生产环境
的正确连接和安全性。
```
2026-05-12 16:07:07 +08:00
7d7c471b39 ```
feat(agents): 更新生产环境意图识别服务的API密钥

- 将INTENT_RECOGNITION服务的app-id从d3c98sp0gdon0fcf6l00更新为d5dn5ibp9adhq1b34lig- 将INTENT_RECOGNITION服务的app-key从d3c98vp0gdon0fcf6n20更新为d6j5bcellh49on5tasvg
```
2026-05-12 15:49:59 +08:00
3868194932 refactor(flow): 优化流程控制异步处理和安全配置
- 将任务停止操作改为异步执行以提高性能
- 替换系统输出为日志记录以改进调试追踪
- 为安全工具类添加异常处理和匿名用户支持
- 更新安全配置允许流程执行接口匿名访问
- 为流程控制器添加匿名注解支持外部调用
2026-05-11 11:47:46 +08:00
0dbed4006c refactor(llm): 重构AI代理平台服务接口和实现
- 修改LLMAiAgentPlatformService.query方法签名,添加TTS相关参数
- 更新FlowActionExecutorService中的AI代理平台调用逻辑
- 移除FlowHttpNodeHandler中的formData清理逻辑
- 简化FlowNodeParamPreparer中参数准备逻辑
- 更新LLMAiAgentPlatformOperateService参数解析结构
- 修改LlmChatService.chat方法添加TTS配置参数
- 实现TTS播放控制功能,支持配置化语音合成参数
- 更新意图识别和工作流调用服务兼容性适配
2026-04-17 09:10:58 +08:00
37d800be69 feat(flow): 添加AI代理平台支持并优化流式处理
- 在FlowActionExecutorService中集成LLMAiAgentPlatformService和LLMAiTtsService
- 实现AI_AGENT_PLATFORM类型的执行逻辑,支持代理平台查询功能
- 重构FlowHttpNodeHandler参数解析,支持JSON和表单数据格式
- 在FlowItemExecutor中清空起始节点输入定义
- 重写FlowNodeParamPreparer参数准备逻辑,支持嵌套参数处理
- 优化LlmChatService中的SSE流处理,实现句子级别异步处理
- 添加TTS语音播放功能,在流式响应中实时处理句子
- 修改评估服务调用方式,在FlowEndNodeHandler中注释掉原有逻辑
2026-04-15 11:50:01 +08:00
722e8a6ba7 Merge branch 'dev2' into dev
# Conflicts:
#	cmvr-iot-test/src/main/java/com/cmvr/test/flow/control/FlowControlService.java
2026-04-09 17:33:11 +08:00
fee5d2efb9 Merge remote-tracking branch 'origin/dev2' into dev2
# Conflicts:
#	cmvr-iot-test/src/main/java/com/cmvr/test/flow/control/FlowControlService.java
2026-04-09 17:30:44 +08:00
stream
3170e86716 refactor(flow): 优化流程控制异步执行和节点查找逻辑
- 将终端停止操作改为异步执行以提升性能
- 替换子图开始节点查找方法为入度为零节点查找
- 扩展子图结束节点判断条件增强流程控制准确性
- 移除异常捕获中的空处理块简化代码结构
2026-04-09 17:29:47 +08:00
d8684a7736 refactor(flow): 优化暂停任务逻辑以异步停止边缘端
- 将边缘端停止操作改为异步执行避免阻塞主线程
- 添加日志记录开始暂停操作
- 保持原有的暂停状态设置和时间戳更新逻辑
- 继续使用统一的日志记录和状态同步机制
2026-04-09 14:42:14 +08:00
stream
a1c3040210 fix(flow): 修复流程控制中的循环数组传递问题
- 在EdgeHlcServiceImpl中添加触控坐标验证,防止坐标为0时的异常
- 在FlowControlService中为resumeMessage添加loopArray属性传递,确保恢复时保留循环数组数据
- 修改FlowGetCurrentObjNodeHandler中获取当前对象的方式,使用loopArray替代inputParams中的array
- 在FlowItemExecutor中将rootMessage的loopArray传递给子消息
- 更新FlowLoopNodeHandler中循环处理逻辑,确保loopArray在循环过程中正确传递
- 调整FlowNodeParamPreparer中循环目标解析方式,直接使用loopNumVal作为目标对象
- 在TaskNodeExecuteContext和TaskNodeExecuteMessage中添加loopArray属性支持
2026-04-09 14:29:31 +08:00
stream
8065faebcb fix: 触控交互 2026-03-24 10:59:50 +08:00
25 changed files with 311 additions and 109 deletions

View File

@ -1,5 +1,6 @@
package com.cmvr.web.controller.test;
import com.cmvr.common.annotation.Anonymous;
import com.cmvr.common.core.controller.BaseController;
import com.cmvr.common.core.domain.AjaxResult;
import com.cmvr.test.flow.context.TaskContextManager;
@ -26,6 +27,8 @@ import org.springframework.web.bind.annotation.RequestParam;
import org.springframework.web.bind.annotation.RestController;
import javax.validation.Valid;
import java.util.ArrayList;
import java.util.List;
@Api(tags = "测试--流程服务")
@RestController
@ -40,15 +43,34 @@ public class TeFlowController extends BaseController {
private final FlowControlService flowControlService;
private final FlowActionExecutorService flowActionExecutorService;
private final TiTouchOperateService tiTouchOperateService;
private final List<String> taskInstIdList = new ArrayList<>();
@ApiOperation("流程发布")
@PostMapping("/publish")
public AjaxResult publish(@RequestBody TeDetectionItem detectionItem) {
return toAjax(teDetectionItemService.publish(detectionItem.getId()));
}
@ApiOperation("任务执行入口")
@PostMapping("/demo")
@Anonymous
public AjaxResult demo(@Valid @RequestBody TeTaskExecuteNormalVO taskExecuteNormalVO) {
if (taskInstIdList.size() > 0) {
try {
taskInstIdList.forEach(flowControlService::stop);
} catch (Exception e) {
// 忽略
}
taskInstIdList.clear();
}
String insId = flowTaskRuntimeService.executeTask(taskExecuteNormalVO);
taskInstIdList.add(insId);
return AjaxResult.success();
}
@ApiOperation("任务执行入口")
@PostMapping("/execute")
@Anonymous
public AjaxResult execute(@Valid @RequestBody TeTaskExecuteNormalVO taskExecuteNormalVO) {
return AjaxResult.success(flowTaskRuntimeService.executeTask(taskExecuteNormalVO));
}

View File

@ -92,7 +92,7 @@ spring:
# Minio配置
minio:
url: http://mini:9000
url: http://minio:9000
accessKey: AKICMVR
secretKey: wJalrXUtnFEMI
bucketName: cmvr-iot
@ -118,8 +118,8 @@ api:
# ActionEnum 的枚举作为 key
agents:
INTENT_RECOGNITION: # 意图识别
app-id: d3c98sp0gdon0fcf6l00
app-key: d3c98vp0gdon0fcf6n20
app-id: d5dn5ibp9adhq1b34lig
app-key: d6j5bcellh49on5tasvg
TI_TOUCH_COORDINATES: # 获取触控二维坐标
app-id: d1ebtabnjkflk4gmhikg
app-key: d5thge2cktmipk78h82g

View File

@ -3,6 +3,7 @@ package com.cmvr.edge.client.service.impl;
import cmvr.api.HlcCommand;
import cmvr.api.HlcServiceGrpc;
import com.alibaba.fastjson2.JSON;
import com.cmvr.common.exception.GlobalException;
import com.cmvr.edge.client.manage.GrpcServiceManager;
import com.cmvr.edge.client.model.hlc.EdgeTouchVO;
import com.cmvr.edge.client.service.EdgeHlcService;
@ -85,6 +86,9 @@ public class EdgeHlcServiceImpl implements EdgeHlcService {
@Override
public String touch(EdgeTouchVO edgeTouchVO) {
if (edgeTouchVO.getX() == edgeTouchVO.getY() && edgeTouchVO.getX() == 0) {
throw new GlobalException("触控坐标不能为0");
}
HlcServiceGrpc.HlcServiceBlockingStub stub = grpcServiceManager.getGrpcClient(edgeTouchVO.getTerminalId(), HlcServiceGrpc.HlcServiceBlockingStub.class);
HlcCommand.Touch.Request request = HlcCommand.Touch.Request.newBuilder()
.setHeader(EdgeCommonUtil.buildRequest(edgeTouchVO.getDeviceId()))

View File

@ -3,5 +3,5 @@ package com.cmvr.llm.service;
import com.alibaba.fastjson2.JSONObject;
public interface LLMAiAgentPlatformService {
JSONObject query(String action, String text, String apiKey);
JSONObject query(String action, String text, String apiKey, Boolean invokeTts, JSONObject config);
}

View File

@ -90,7 +90,9 @@ public class WorkflowInvokeServiceCopy {
appKey, // apiKey (header用)
appId, // apiId (body的AppKey)
"user_123", // userId
body
body,
false,
null
);
System.out.println( result);
System.out.println("总耗时:" + (System.currentTimeMillis() - startTime));

View File

@ -14,7 +14,7 @@ import java.util.Map;
public class LLMAiAgentPlatformServiceImpl implements LLMAiAgentPlatformService {
private final LlmChatService llmChatService;
@Override
public JSONObject query(String action, String text, String apiKey) {
public JSONObject query(String action, String text, String apiKey, Boolean invokeTts, JSONObject config) {
Map<String, Object> body = new HashMap<>();
body.put("Query", text);
long startTime = System.currentTimeMillis();
@ -24,7 +24,9 @@ public class LLMAiAgentPlatformServiceImpl implements LLMAiAgentPlatformService
apiKey, // apiKey (header用)
apiKey, // apiId (body的AppKey)
"user_123", // userId
body
body,
invokeTts,
config
);
System.out.println( result);
System.out.println("总耗时:" + (System.currentTimeMillis() - startTime));

View File

@ -36,7 +36,9 @@ public class LLMIntentRecognitionServiceImpl implements LLMIntentRecognitionServ
acg.getAppKey(), // apiKey (header用)
acg.getAppId(), // apiId (body的AppKey)
"user_123", // userId
body
body,
false,
null
);
System.out.println( result);
System.out.println("总耗时:" + (System.currentTimeMillis() - startTime));

View File

@ -1,7 +1,11 @@
package com.cmvr.llm.util;
import com.alibaba.fastjson2.JSONObject;
import com.cmvr.llm.service.LLMAiTtsService;
import com.fasterxml.jackson.databind.ObjectMapper;
import lombok.extern.slf4j.Slf4j;
import okhttp3.*;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.stereotype.Service;
import javax.annotation.PreDestroy;
@ -10,17 +14,18 @@ import java.io.IOException;
import java.io.InputStreamReader;
import java.nio.charset.StandardCharsets;
import java.util.*;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.CountDownLatch;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.*;
import java.util.concurrent.atomic.AtomicBoolean;
import java.util.concurrent.atomic.AtomicReference;
/**
* LLM聊天服务 - 严格保持原接口参数
*/
@Slf4j
@Service
public class LlmChatService {
@Autowired
private LLMAiTtsService llmAiTtsService;
private static final String DEFAULT_BASE_URL = "https://aiagentplatform.cmft.com/api/proxy/api/v1";
private static final MediaType JSON = MediaType.parse("application/json; charset=utf-8");
private static final ObjectMapper objectMapper = new ObjectMapper();
@ -48,7 +53,7 @@ public class LlmChatService {
* @param body 请求体必须包含Query字段可选Name等
* @return 完整AI回复字符串
*/
public String chat(String apiKey, String apiId, String userId, Map<String, Object> body) throws IOException {
public String chat(String apiKey, String apiId, String userId, Map<String, Object> body, boolean playTts, JSONObject ttsConfig) throws IOException {
String cacheKey = apiKey + "|" + apiId + "|" + userId;
// 获取或创建会话严格保持原createConversation逻辑
@ -76,7 +81,7 @@ public class LlmChatService {
}
// 执行流式请求并同步返回完整结果
return executeSseRequestSync(apiKey, requestBody);
return executeSseRequestSync(apiKey, requestBody, playTts, ttsConfig);
}
/**
@ -122,10 +127,13 @@ public class LlmChatService {
}
}
// 句子结束符号
static final String SENTENCE_END = "?!。?!\n";
/**
* 执行SSE请求 - 完全保持原executeSseRequest逻辑改为同步返回
*/
private String executeSseRequestSync(String apiKey, Map<String, Object> requestBody) throws IOException {
private String executeSseRequestSync(String apiKey, Map<String, Object> requestBody, boolean playTts, JSONObject ttsConfig) throws IOException {
String jsonBody = objectMapper.writeValueAsString(requestBody);
Request request = new Request.Builder()
@ -137,18 +145,49 @@ public class LlmChatService {
.build();
StringBuilder fullContent = new StringBuilder();
CountDownLatch latch = new CountDownLatch(1);
CountDownLatch mainLatch = new CountDownLatch(1);
AtomicReference<Throwable> errorRef = new AtomicReference<>();
AtomicReference<String> finalMessageIdRef = new AtomicReference<>("");
// ====================== 核心控制 ======================
// 接收缓存永远不阻塞
StringBuilder receiveBuffer = new StringBuilder();
// 单线程串行执行保证上一句处理完才处理下一句
ExecutorService sentenceExecutor = Executors.newSingleThreadExecutor();
// 标记流是否已经结束
AtomicBoolean streamFinished = new AtomicBoolean(false);
final AtomicBoolean interrupted = new AtomicBoolean(false); // 中断标记
Call call = httpClient.newCall(request);
Runnable interruptTask = () -> {
interrupted.set(true);
call.cancel(); // 真正关闭SSE连接 关键
streamFinished.set(true);
sentenceExecutor.shutdownNow();
mainLatch.countDown();
};
// 使用异步调用但阻塞等待结果
call.enqueue(new Callback() {
/**
* 等待所有句子处理完再结束
*/
private void shutdownAndWait1() {
try {
sentenceExecutor.shutdown();
sentenceExecutor.awaitTermination(10, TimeUnit.MINUTES);
} catch (InterruptedException ignored) {}
mainLatch.countDown();
}
@Override
public void onFailure(Call call, IOException e) {
errorRef.set(e);
latch.countDown();
streamFinished.set(true);
shutdownAndWait1();
}
@Override
@ -156,7 +195,8 @@ public class LlmChatService {
if (!response.isSuccessful()) {
String errorBody = response.body() != null ? response.body().string() : "无错误信息";
errorRef.set(new IOException("SSE请求失败: " + response.code() + " - " + errorBody));
latch.countDown();
streamFinished.set(true);
shutdownAndWait1();
return;
}
@ -168,7 +208,7 @@ public class LlmChatService {
new InputStreamReader(response.body().byteStream(), StandardCharsets.UTF_8))) {
String line;
while ((line = reader.readLine()) != null) {
while (!interrupted.get() && (line = reader.readLine()) != null) {
if (line.isEmpty()) {
currentEvent = "";
continue;
@ -190,8 +230,9 @@ public class LlmChatService {
if ("[DONE]".equals(data)) {
fullContent.append(localContent);
finalMessageIdRef.set(finalMessageId);
latch.countDown();
return;
streamFinished.set(true);
processBufferIfNeed();
break;
}
try {
@ -208,15 +249,18 @@ public class LlmChatService {
String answer = extractString(jsonData, "answer", "content", "chunk", "text");
if (!answer.isEmpty()) {
localContent.append(answer);
receiveBuffer.append(answer);
processBufferIfNeed();
}
} else if ("message_start".equals(event)) {
System.out.println("消息开始ID: " + finalMessageId);
log.info("消息开始ID: " + finalMessageId);
} else if ("message_output_start".equals(event)) {
System.out.println("输出开始");
log.info("输出开始");
} else if ("message_end".equals(event) || "end".equals(event)) {
fullContent.append(localContent);
finalMessageIdRef.set(finalMessageId);
latch.countDown();
streamFinished.set(true);
processBufferIfNeed();
return;
}
@ -229,22 +273,77 @@ public class LlmChatService {
// 流正常结束
fullContent.append(localContent);
finalMessageIdRef.set(finalMessageId);
latch.countDown();
} catch (IOException e) {
errorRef.set(e);
latch.countDown();
streamFinished.set(true);
} finally {
processBufferIfNeed();
shutdownAndWait1();
}
}
/**
* 核心提取完整句子 + 异步串行处理
*/
private void processBufferIfNeed() {
if (interrupted.get()) return;
if (!true) return;
if (!playTts) return;
sentenceExecutor.submit(() -> {
while (true) {
String buffer = receiveBuffer.toString();
// 寻找最后一个句子结束符
int lastSplit = -1;
for (int i = 0; i < buffer.length(); i++) {
if (SENTENCE_END.contains(String.valueOf(buffer.charAt(i)))) {
lastSplit = i;
}
}
// 有完整句子 流已结束
if (lastSplit >= 0 || streamFinished.get()) {
String sentence;
if (lastSplit >= 0) {
// 截取完整句子
sentence = buffer.substring(0, lastSplit + 1).trim();
// 保留剩余内容
String remain = buffer.substring(lastSplit + 1).trim();
receiveBuffer.setLength(0);
receiveBuffer.append(remain);
} else {
// 流结束直接处理剩余
sentence = buffer.trim();
receiveBuffer.setLength(0);
}
if (!sentence.isEmpty()) {
if (interrupted.get()) return;
yourAsyncMethod(sentence, ttsConfig);
}
// 处理完继续循环看是否还有新句子
if (streamFinished.get() && receiveBuffer.length() == 0) {
break;
}
} else {
// 无完整句子退出
break;
}
}
});
}
});
// 等待流式传输完成
try {
boolean completed = latch.await(120, TimeUnit.SECONDS);
boolean completed = mainLatch.await(120, TimeUnit.SECONDS);
if (!completed) {
throw new IOException("请求超时");
}
} catch (InterruptedException e) {
interruptTask.run(); // 被中断时立即停止SSE
Thread.currentThread().interrupt();
throw new IOException("请求被中断", e);
}
@ -260,6 +359,15 @@ public class LlmChatService {
return fullContent.toString();
}
private void yourAsyncMethod(String sentence, JSONObject ttsConfig) {
try {
log.info("正在处理句子: " + sentence);
llmAiTtsService.play(ttsConfig.getString("url"), sentence, ttsConfig.getString("voice"), ttsConfig.getString("speed"), ttsConfig.getString("volume"));
} catch (Exception e) {
throw new RuntimeException(e);
}
}
/**
* 完全保持原extractString方法
*/

View File

@ -115,7 +115,7 @@ public class SecurityConfig
// 静态资源可匿名访问
.antMatchers(HttpMethod.GET, "/", "/*.html", "/**/*.html", "/**/*.css", "/**/*.js", "/profile/**").permitAll()
.antMatchers("/system/file/upload","/evaluation/callback","/flow/**","/flowise/**","/kws/**","/ws/**","/api/grpc/**","/show/**",
"/node-red/**","/swagger-ui.html", "/swagger-resources/**", "/webjars/**", "/*/api-docs", "/druid/**", "/ti/**").permitAll()
"/node-red/**","/swagger-ui.html", "/swagger-resources/**", "/webjars/**", "/*/api-docs", "/druid/**", "/ti/**", "/flow/execute").permitAll()
// 除上面外的所有请求全部需要鉴权认证
.anyRequest().authenticated();
})

View File

@ -24,10 +24,17 @@ public class MyMetaObjectHandler implements MetaObjectHandler {
*/
@Override
public void insertFill(MetaObject metaObject) {
// 获取用户名异常时给默认值 anonymous
String username;
try {
username = SecurityUtils.getUsername();
} catch (Exception e) {
username = "anonymous"; // 匿名接口默认值
}
// 起始版本 3.3.0(推荐使用)
this.setFieldValByName(CREATE_BY, SecurityUtils.getUsername(), metaObject);
this.setFieldValByName(CREATE_BY, username, metaObject);
this.setFieldValByName(CREATE_TIME, formatDate(metaObject.getSetterType(CREATE_TIME)), metaObject);
this.setFieldValByName(UPDATE_BY, SecurityUtils.getUsername(), metaObject);
this.setFieldValByName(UPDATE_BY, username, metaObject);
this.setFieldValByName(UPDATE_TIME, formatDate(metaObject.getSetterType(CREATE_TIME)), metaObject);
this.setFieldValByName(DELETED, "0", metaObject);
// this.setFieldValByName(STATUS, "1", metaObject);
@ -41,7 +48,14 @@ public class MyMetaObjectHandler implements MetaObjectHandler {
*/
@Override
public void updateFill(MetaObject metaObject) {
this.setFieldValByName(UPDATE_BY, SecurityUtils.getUsername(), metaObject);
// 获取用户名异常时给默认值 anonymous
String username;
try {
username = SecurityUtils.getUsername();
} catch (Exception e) {
username = "anonymous"; // 匿名接口默认值
}
this.setFieldValByName(UPDATE_BY, username, metaObject);
this.setFieldValByName(UPDATE_TIME, formatDate(metaObject.getSetterType(CREATE_TIME)), metaObject);
}

View File

@ -54,7 +54,10 @@ public class FlowControlService {
if (ctx.isStopped()) {
throw new GlobalException("任务已终止,无需重复操作");
}
edgeSystemService.stopAll(ctx.getTerminalId());
// 异步终止
executor.execute(() -> {
edgeSystemService.stopAll(ctx.getTerminalId());
});
// 统一记录日志 + 设置上下文状态 + 数据库状态
taskInstHolder.syncStatus(instId, TaskStatusEnum.STOPPED);
@ -78,7 +81,10 @@ public class FlowControlService {
if (ctx.isPaused()) {
throw new GlobalException("任务已处于暂停状态");
}
edgeSystemService.stopAll(ctx.getTerminalId());
// 异步停止终端
executor.execute(() -> {
edgeSystemService.stopAll(ctx.getTerminalId());
});
ctx.setPaused(true);
ctx.setTerminalStatus(TerminalStatusEnum.RUNNING);
ctx.setStatus(TaskStatusEnum.PAUSED);
@ -144,6 +150,7 @@ public class FlowControlService {
TaskNodeExecuteMessage resumeMessage = new TaskNodeExecuteMessage();
BeanUtil.copyProperties(pendingNode.getRootMessage(), resumeMessage);
resumeMessage.setLoopNum(pendingNode.getLoopIteration()); // 保留中断时的循环次数
resumeMessage.setLoopArray(pendingNode.getLoopArray()); // 保留中断时的循环次数
resumeMessage.setIterations(new ArrayList<>(pendingNode.getIterations())); // 保留中断时的路径
flowItemExecutor.executeNode(
@ -169,6 +176,7 @@ public class FlowControlService {
TaskNodeExecuteMessage resumeMessage = new TaskNodeExecuteMessage();
BeanUtil.copyProperties(pendingNode.getRootMessage(), resumeMessage);
resumeMessage.setLoopNum(pendingNode.getLoopIteration()); // 保留循环次数
resumeMessage.setLoopArray(pendingNode.getLoopArray()); // 保留循环次数
resumeMessage.setIterations(new ArrayList<>(pendingNode.getIterations())); // 保留路径
flowItemExecutor.executeNode(

View File

@ -47,7 +47,7 @@ public class FlowEndNodeHandler implements FlowNodeTypeHandler {
Thread.currentThread().setName("evaluation-thread-" + Thread.currentThread().getId());
try {
log.info("异步评估任务开始执行threadName={}, instId={}, itemId={}", Thread.currentThread().getName(), instId, itemId);
exAeEvaluationService.executeEvaluation(jsonObject);
// exAeEvaluationService.executeEvaluation(jsonObject);
log.info("异步评估任务执行结束threadName={}, instId={}, itemId={}", Thread.currentThread().getName(), instId, itemId);
} catch (Exception e) {
log.error("评估任务执行异常threadName={}, instId={}, itemId={}, 错误={}", Thread.currentThread().getName(), instId, itemId, e.getMessage(), e);

View File

@ -8,6 +8,7 @@ import com.cmvr.test.flow.runtime.message.TaskNodeExecuteResult;
import lombok.extern.slf4j.Slf4j;
import org.springframework.stereotype.Component;
import java.util.Collections;
import java.util.List;
@Slf4j
@ -18,17 +19,17 @@ public class FlowGetCurrentObjNodeHandler implements FlowNodeTypeHandler {
@Override
public TaskNodeExecuteResult handle(TaskNodeExecuteMessage message) {
try {
JSONObject inputParams = message.getInputParams();
JSONArray jsonArray = inputParams.getJSONArray("array");
List<Integer> iterations = message.getIterations();
int last = CollUtil.getLast(iterations) - 1;
// JSONObject inputParams = message.getInputParams();
// JSONArray jsonArray = inputParams.getJSONArray("array");
//
// List<Integer> iterations = message.getIterations();
// int last = CollUtil.getLast(iterations) - 1;
int index = message.getLoopNum() - 1;
JSONObject output = new JSONObject();
output.put("index", last);
if (CollUtil.isNotEmpty(jsonArray) && last >= 0) {
output.put("object", jsonArray.get(last));
output.put("index", index);
JSONArray objects = (JSONArray)message.getLoopArray();
if (CollUtil.isNotEmpty(objects) && index >= 0) {
output.put("object", objects.get(index));
}
return TaskNodeExecuteResult.success(output);

View File

@ -20,22 +20,30 @@ public class FlowHttpNodeHandler implements FlowNodeTypeHandler {
public TaskNodeExecuteResult handle(TaskNodeExecuteMessage message) {
try {
JSONObject inputParams = message.getInputParams();
JSONObject config = inputParams.getJSONObject("config");
JSONObject body = inputParams.getJSONObject("body");
JSONObject headers = inputParams.getJSONObject("headers");
// 1. 必填参数校验
String url = inputParams.getString("url");
String url = config.getString("url");
if (url == null || url.isEmpty()) {
return TaskNodeExecuteResult.failure("HTTP url is required");
}
String method = inputParams.getString("method");
String method = config.getString("method");
if (method == null) {
method = "POST";
}
int timeout = inputParams.getIntValue("timeout", 100000);
int timeout = config.getIntValue("timeout", 100000);
JSONObject headers = inputParams.getJSONObject("headers");
Object bodyCfg = inputParams.get("body");
Object bodyCfg = null;
if (!"json".equals(body.getString("bodyType"))) {
JSONObject formData = body.getJSONObject("formData");
bodyCfg = formData;
} else {
bodyCfg = body.get("json");
}
// 2. 构造 HTTP 请求
HttpRequest request = HttpRequest.of(url)
.method(Method.valueOf(method.toUpperCase()))

View File

@ -37,6 +37,7 @@ public class FlowLoopNodeHandler implements FlowNodeTypeHandler {
String nodeId = message.getNodeId();
FlowNodeWrapper nodeWrapper = graph.getNode(nodeId);
int loopCount = inputParams.getIntValue("loopNum"); // 获取循环次数
Object loopArray = inputParams.get("loopArray"); // 获取循环次数
// 如果没有设置循环次数或者循环次数为 0则跳过处理
if (loopCount <= 0) {
@ -65,8 +66,9 @@ public class FlowLoopNodeHandler implements FlowNodeTypeHandler {
// 克隆 message并明确设置 loopNum
TaskNodeExecuteMessage subMessage = new TaskNodeExecuteMessage();
BeanUtil.copyProperties(message, subMessage);
// subMessage.setLoopNum(i); // 当前 loop i
subMessage.setLoopNum(i); // 当前 loop i
subMessage.setIterations(newIterations); // 完整路径
subMessage.setLoopArray(loopArray);
flowItemExecutor.executeSubGraph(subGraph, subMessage, latch::countDown, newIterations);

View File

@ -50,7 +50,7 @@ public class FlowItemExecutor {
FlowGraph graph = FlowModelBuilder.buildExecutableGraph(item.getFlowData());
String startNodeId = graph.findStartNodeId();
List<FlowParamDef> startInputDefs = graph.getStartInputParamsDefs();
startInputDefs.clear();
// 执行 start 节点同步执行
NodeExecutor executor = node ->
executeNode(graph, node, startNodeId, startInputDefs, rootMessage, onFinished, new ArrayList<>());
@ -73,7 +73,7 @@ public class FlowItemExecutor {
List<FlowParamDef> startInputDefs = subGraph.getStartInputParamsDefs();
// 子图的 start 节点
FlowNodeWrapper subStartNode = subGraph.findSubStartNodeId();
FlowNodeWrapper subStartNode = subGraph.getInDegreeZeroNode();
String subStartNodeId = subStartNode.getNodeId();
NodeExecutor executor = node -> executeNode(
@ -84,7 +84,7 @@ public class FlowItemExecutor {
// 执行子图起始节点
executeNode(subGraph, subStartNode, subStartNodeId, startInputDefs,
rootMessage, onFinished, iterations);
log.info("当前执行-子图");
// 启动子图调度
flowTaskScheduler.subStart(subGraph, taskInstHolder.getContext(rootMessage.getInstId()),
executor, onFinished, iterations);
@ -131,7 +131,7 @@ public class FlowItemExecutor {
}
// END 节点 or 子图结束
if (node.getNodeType() == NodeTypeEnum.END || node.getNodeType() == NodeTypeEnum.SUB_END) {
if (node.getNodeType() == NodeTypeEnum.END || node.getNodeType() == NodeTypeEnum.SUB_END || graph.isSubGraphEndNode(nodeId)) {
try {
log.info("{} 节点执行,释放线程", node.getNodeType());
onFinished.run();
@ -172,6 +172,8 @@ public class FlowItemExecutor {
// loopNum 不在这里计算而是由 FlowLoopNodeHandler / resume 显式写入
// 如果 rootMessage 里已经带了 loopNum就沿用它
message.setLoopNum(rootMessage.getLoopNum());
message.setLoopArray(rootMessage.getLoopArray());
return message;
}

View File

@ -42,6 +42,8 @@ public class FlowTaskRuntimeEntry implements FlowTaskRuntimeService {
private final TaskInstHolder taskInstHolder;
private final FlowTaskAsyncDispatcher flowTaskAsyncDispatcher;
@Transactional(rollbackFor = Exception.class)
public String executeTask(TeTaskExecuteNormalVO taskExecuteNormalVO) {
return executeTaskInternal(

View File

@ -18,6 +18,7 @@ public interface FlowTaskRuntimeService {
*/
String executeTask(TeTaskExecuteNormalVO taskExecuteNormalVO);
/**
* 执行试运行任务TRIAL 模式
*

View File

@ -58,7 +58,7 @@ public class FlowTaskScheduler {
Runnable onFinished,
List<Integer> iterations) {
String instId = context.getInstId();
FlowNodeWrapper startNode = graph.findSubStartNodeId();
FlowNodeWrapper startNode = graph.getInDegreeZeroNode();
List<String> nextNodes = graph.getNextNodes(startNode.getNodeId());
if (nextNodes.isEmpty()) {

View File

@ -1,6 +1,7 @@
package com.cmvr.test.flow.runtime.engine.support;
import cn.hutool.core.collection.CollUtil;
import cn.hutool.core.collection.CollectionUtil;
import com.alibaba.fastjson2.JSON;
import com.alibaba.fastjson2.JSONArray;
import com.alibaba.fastjson2.JSONObject;
@ -30,11 +31,65 @@ public class FlowNodeParamPreparer {
FlowNodeWrapper node,
TaskNodeExecuteMessage rootMessage) {
JSONObject input = new JSONObject();
TaskContext ctx = taskInstHolder.getContext(rootMessage.getInstId());
JSONObject input = getInputParams(graph, node.getNodeParams(), rootMessage);
// 3 START 节点 runParams 覆盖默认参数
if (node.getNodeType() == NodeTypeEnum.START) {
JSONObject merged = new JSONObject(input);
JSONObject runParams = ctx.getRunParams();
for (String key : runParams.keySet()) {
// runParams 里有值就覆盖掉定义里的值
merged.put(key, runParams.get(key));
}
return merged;
}
// 4 LOOP 节点动态计算 loopCount
if (node.getNodeType() == NodeTypeEnum.LOOP) {
Object loopNumVal = input.get("loopNum");
input.put("loopArray", loopNumVal);
int loopCount = 0;
// 根据迭代路径找到当前层的集合对象
// Object target = resolveLoopTarget(loopNumVal, rootMessage.getIterations(), 0);
Object target = loopNumVal;
if (target instanceof JSONArray) {
loopCount = ((JSONArray) target).size();
} else if (target instanceof Collection) {
loopCount = ((Collection<?>) target).size();
} else if (target != null && target.getClass().isArray()) {
loopCount = Array.getLength(target);
} else if (target instanceof Number) {
loopCount = ((Number) target).intValue();
} else if (target != null) {
try {
String string = target.toString();
if (string.startsWith("[")) {
JSONArray array = JSON.parseArray(string);
loopCount = array.size();
}else {
loopCount = Integer.parseInt(string);
}
} catch (NumberFormatException e) {
throw new GlobalException("loopNum 参数格式错误: " + loopNumVal);
}
}
// 写回 inputParams LoopHandler 使用
input.put("loopNum", loopCount);
}
return input;
}
private JSONObject getInputParams(FlowGraph graph, List<FlowParamDef> paramDefList, TaskNodeExecuteMessage rootMessage) {
JSONObject input = new JSONObject();
TaskContext ctx = taskInstHolder.getContext(rootMessage.getInstId());
// 遍历当前节点定义的参数
for (FlowParamDef param : node.getNodeParams()) {
for (FlowParamDef param : paramDefList) {
if (param.getScope() != ParamScope.NODE) continue;
// 1 静态 input 参数
@ -68,53 +123,11 @@ public class FlowNodeParamPreparer {
Object val = findNestedValue(source, path);
input.put(param.getName(), val);
}
}
// 3 START 节点 runParams 覆盖默认参数
if (node.getNodeType() == NodeTypeEnum.START) {
JSONObject merged = new JSONObject(input);
JSONObject runParams = ctx.getRunParams();
for (String key : runParams.keySet()) {
// runParams 里有值就覆盖掉定义里的值
merged.put(key, runParams.get(key));
if (!CollectionUtil.isEmpty(param.getChildren())) {
// 添加子参数
input.put(param.getName(), getInputParams(graph, param.getChildren(), rootMessage));
}
return merged;
}
// 4 LOOP 节点动态计算 loopCount
if (node.getNodeType() == NodeTypeEnum.LOOP) {
Object loopNumVal = input.get("loopNum");
int loopCount = 0;
// 根据迭代路径找到当前层的集合对象
Object target = resolveLoopTarget(loopNumVal, rootMessage.getIterations(), 0);
if (target instanceof JSONArray) {
loopCount = ((JSONArray) target).size();
} else if (target instanceof Collection) {
loopCount = ((Collection<?>) target).size();
} else if (target != null && target.getClass().isArray()) {
loopCount = Array.getLength(target);
} else if (target instanceof Number) {
loopCount = ((Number) target).intValue();
} else if (target != null) {
try {
String string = target.toString();
if (string.startsWith("[")) {
JSONArray array = JSON.parseArray(string);
loopCount = array.size();
}else {
loopCount = Integer.parseInt(string);
}
} catch (NumberFormatException e) {
throw new GlobalException("loopNum 参数格式错误: " + loopNumVal);
}
}
// 写回 inputParams LoopHandler 使用
input.put("loopNum", loopCount);
}
return input;
}

View File

@ -22,5 +22,6 @@ public class TaskNodeExecuteContext {
private TaskNodeExecuteMessage rootMessage;
private Runnable onFinished;
private int loopIteration;
private Object loopArray;
private List<Integer> iterations = new ArrayList<>();
}

View File

@ -55,6 +55,8 @@ public class TaskNodeExecuteMessage {
*/
private List<Integer> iterations = new ArrayList<>();
private Object loopArray;
/**
* 节点名称
*/

View File

@ -12,6 +12,7 @@ import org.springframework.stereotype.Service;
@Slf4j
@Service
@RequiredArgsConstructor
public class LLMAiAgentPlatformOperateService implements LLMOperateService {
@ -26,10 +27,11 @@ public class LLMAiAgentPlatformOperateService implements LLMOperateService {
public TaskNodeExecuteResult execute(TaskNodeExecuteMessage message) {
ActionEnum action = message.getAction();
JSONObject inputParams = message.getInputParams();
String text = inputParams.getString("text");
String apiKey = inputParams.getString("apiKey");
JSONObject output = llmAiAgentPlatformService.query(action.name(), text, apiKey);
JSONObject config = inputParams.getJSONObject("config");
String text = config.getString("text");
String apiKey = config.getString("apiKey");
Boolean invokeTts = inputParams.getBoolean("invokeTts");
JSONObject output =llmAiAgentPlatformService.query(action.name(), text, apiKey, invokeTts, inputParams.getJSONObject("tts"));
return TaskNodeExecuteResult.success(output);
}

View File

@ -12,6 +12,8 @@ import com.cmvr.edge.client.service.EdgeCameraService;
import com.cmvr.edge.client.service.EdgeHlcService;
import com.cmvr.edge.client.service.EdgeMicrophoneService;
import com.cmvr.edge.client.service.EdgeSpeakerService;
import com.cmvr.llm.service.LLMAiAgentPlatformService;
import com.cmvr.llm.service.LLMAiTtsService;
import com.cmvr.test.enums.ActionEnum;
import com.cmvr.test.model.vo.FlowActionRequestVO;
import lombok.RequiredArgsConstructor;
@ -27,6 +29,8 @@ public class FlowActionExecutorService {
private final EdgeMicrophoneService edgeMicrophoneService;
private final EdgeSpeakerService edgeSpeakerService;
private final EdgeHlcService edgeHlcService;
private final LLMAiTtsService llmAiTtsService;
private final LLMAiAgentPlatformService llmAiAgentPlatformService;
public String actionExecute(FlowActionRequestVO req) {
@ -115,6 +119,8 @@ public class FlowActionExecutorService {
return "LLM 意图识别 OK";
case GENERATE_ADVANCED_AUDIO:
return "LLM 高级音频生成 OK";
case AI_AGENT_PLATFORM:
return llmAiAgentPlatformService.query(action.getAction(), req.getPayload().getJSONObject("config").getString("text"), req.getPayload().getJSONObject("config").getString("apiKey"), req.getPayload().getBoolean("invokeTts"), req.getPayload().getJSONObject("tts")).toString();
default:
throw new UnsupportedOperationException("未实现的 LLM Action: " + action);
}

View File

@ -40,7 +40,7 @@ public class TiVehicleFunctionServiceImpl extends ServiceImpl<TiVehicleFunctionM
@Override
public JSONObject selectTiVehicleFunctionDetailById(Long id) {
JSONObject jsonObject = JSONObject.from(this.getById(id));
jsonObject.put("funcName", id== 26 ? "首页" : "设置");
jsonObject.put("funcName", sysDictDataService.selectDictLabel("ti_function_config", jsonObject.getString("funcKey")));
return jsonObject;
}