Merge branch 'dev2' into dev
# Conflicts: # cmvr-iot-test/src/main/java/com/cmvr/test/flow/control/FlowControlService.java
This commit is contained in:
commit
722e8a6ba7
@ -3,6 +3,7 @@ package com.cmvr.edge.client.service.impl;
|
|||||||
import cmvr.api.HlcCommand;
|
import cmvr.api.HlcCommand;
|
||||||
import cmvr.api.HlcServiceGrpc;
|
import cmvr.api.HlcServiceGrpc;
|
||||||
import com.alibaba.fastjson2.JSON;
|
import com.alibaba.fastjson2.JSON;
|
||||||
|
import com.cmvr.common.exception.GlobalException;
|
||||||
import com.cmvr.edge.client.manage.GrpcServiceManager;
|
import com.cmvr.edge.client.manage.GrpcServiceManager;
|
||||||
import com.cmvr.edge.client.model.hlc.EdgeTouchVO;
|
import com.cmvr.edge.client.model.hlc.EdgeTouchVO;
|
||||||
import com.cmvr.edge.client.service.EdgeHlcService;
|
import com.cmvr.edge.client.service.EdgeHlcService;
|
||||||
@ -85,6 +86,9 @@ public class EdgeHlcServiceImpl implements EdgeHlcService {
|
|||||||
|
|
||||||
@Override
|
@Override
|
||||||
public String touch(EdgeTouchVO edgeTouchVO) {
|
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);
|
HlcServiceGrpc.HlcServiceBlockingStub stub = grpcServiceManager.getGrpcClient(edgeTouchVO.getTerminalId(), HlcServiceGrpc.HlcServiceBlockingStub.class);
|
||||||
HlcCommand.Touch.Request request = HlcCommand.Touch.Request.newBuilder()
|
HlcCommand.Touch.Request request = HlcCommand.Touch.Request.newBuilder()
|
||||||
.setHeader(EdgeCommonUtil.buildRequest(edgeTouchVO.getDeviceId()))
|
.setHeader(EdgeCommonUtil.buildRequest(edgeTouchVO.getDeviceId()))
|
||||||
|
|||||||
@ -78,16 +78,14 @@ public class FlowControlService {
|
|||||||
if (ctx.isPaused()) {
|
if (ctx.isPaused()) {
|
||||||
throw new GlobalException("任务已处于暂停状态");
|
throw new GlobalException("任务已处于暂停状态");
|
||||||
}
|
}
|
||||||
|
// 异步停止终端
|
||||||
log.info("开始暂停");
|
executor.execute(() -> {
|
||||||
|
edgeSystemService.stopAll(ctx.getTerminalId());
|
||||||
|
});
|
||||||
ctx.setPaused(true);
|
ctx.setPaused(true);
|
||||||
ctx.setTerminalStatus(TerminalStatusEnum.RUNNING);
|
ctx.setTerminalStatus(TerminalStatusEnum.RUNNING);
|
||||||
ctx.setStatus(TaskStatusEnum.PAUSED);
|
ctx.setStatus(TaskStatusEnum.PAUSED);
|
||||||
ctx.updateTimestamp();
|
ctx.updateTimestamp();
|
||||||
// 异步停止所有边缘端
|
|
||||||
executor.execute(() -> {
|
|
||||||
edgeSystemService.stopAll(ctx.getTerminalId());
|
|
||||||
});
|
|
||||||
// 统一记录日志 + 设置上下文状态 + 数据库状态
|
// 统一记录日志 + 设置上下文状态 + 数据库状态
|
||||||
taskInstHolder.syncStatus(instId, TaskStatusEnum.PAUSED);
|
taskInstHolder.syncStatus(instId, TaskStatusEnum.PAUSED);
|
||||||
taskThreadRegistry.interruptAll(instId);
|
taskThreadRegistry.interruptAll(instId);
|
||||||
@ -149,6 +147,7 @@ public class FlowControlService {
|
|||||||
TaskNodeExecuteMessage resumeMessage = new TaskNodeExecuteMessage();
|
TaskNodeExecuteMessage resumeMessage = new TaskNodeExecuteMessage();
|
||||||
BeanUtil.copyProperties(pendingNode.getRootMessage(), resumeMessage);
|
BeanUtil.copyProperties(pendingNode.getRootMessage(), resumeMessage);
|
||||||
resumeMessage.setLoopNum(pendingNode.getLoopIteration()); // 保留中断时的循环次数
|
resumeMessage.setLoopNum(pendingNode.getLoopIteration()); // 保留中断时的循环次数
|
||||||
|
resumeMessage.setLoopArray(pendingNode.getLoopArray()); // 保留中断时的循环次数
|
||||||
resumeMessage.setIterations(new ArrayList<>(pendingNode.getIterations())); // 保留中断时的路径
|
resumeMessage.setIterations(new ArrayList<>(pendingNode.getIterations())); // 保留中断时的路径
|
||||||
|
|
||||||
flowItemExecutor.executeNode(
|
flowItemExecutor.executeNode(
|
||||||
@ -174,6 +173,7 @@ public class FlowControlService {
|
|||||||
TaskNodeExecuteMessage resumeMessage = new TaskNodeExecuteMessage();
|
TaskNodeExecuteMessage resumeMessage = new TaskNodeExecuteMessage();
|
||||||
BeanUtil.copyProperties(pendingNode.getRootMessage(), resumeMessage);
|
BeanUtil.copyProperties(pendingNode.getRootMessage(), resumeMessage);
|
||||||
resumeMessage.setLoopNum(pendingNode.getLoopIteration()); // 保留循环次数
|
resumeMessage.setLoopNum(pendingNode.getLoopIteration()); // 保留循环次数
|
||||||
|
resumeMessage.setLoopArray(pendingNode.getLoopArray()); // 保留循环次数
|
||||||
resumeMessage.setIterations(new ArrayList<>(pendingNode.getIterations())); // 保留路径
|
resumeMessage.setIterations(new ArrayList<>(pendingNode.getIterations())); // 保留路径
|
||||||
|
|
||||||
flowItemExecutor.executeNode(
|
flowItemExecutor.executeNode(
|
||||||
|
|||||||
@ -8,6 +8,7 @@ import com.cmvr.test.flow.runtime.message.TaskNodeExecuteResult;
|
|||||||
import lombok.extern.slf4j.Slf4j;
|
import lombok.extern.slf4j.Slf4j;
|
||||||
import org.springframework.stereotype.Component;
|
import org.springframework.stereotype.Component;
|
||||||
|
|
||||||
|
import java.util.Collections;
|
||||||
import java.util.List;
|
import java.util.List;
|
||||||
|
|
||||||
@Slf4j
|
@Slf4j
|
||||||
@ -18,17 +19,17 @@ public class FlowGetCurrentObjNodeHandler implements FlowNodeTypeHandler {
|
|||||||
@Override
|
@Override
|
||||||
public TaskNodeExecuteResult handle(TaskNodeExecuteMessage message) {
|
public TaskNodeExecuteResult handle(TaskNodeExecuteMessage message) {
|
||||||
try {
|
try {
|
||||||
JSONObject inputParams = message.getInputParams();
|
// JSONObject inputParams = message.getInputParams();
|
||||||
JSONArray jsonArray = inputParams.getJSONArray("array");
|
// JSONArray jsonArray = inputParams.getJSONArray("array");
|
||||||
|
//
|
||||||
List<Integer> iterations = message.getIterations();
|
// List<Integer> iterations = message.getIterations();
|
||||||
int last = CollUtil.getLast(iterations) - 1;
|
// int last = CollUtil.getLast(iterations) - 1;
|
||||||
|
int index = message.getLoopNum() - 1;
|
||||||
JSONObject output = new JSONObject();
|
JSONObject output = new JSONObject();
|
||||||
output.put("index", last);
|
output.put("index", index);
|
||||||
|
JSONArray objects = (JSONArray)message.getLoopArray();
|
||||||
if (CollUtil.isNotEmpty(jsonArray) && last >= 0) {
|
if (CollUtil.isNotEmpty(objects) && index >= 0) {
|
||||||
output.put("object", jsonArray.get(last));
|
output.put("object", objects.get(index));
|
||||||
}
|
}
|
||||||
|
|
||||||
return TaskNodeExecuteResult.success(output);
|
return TaskNodeExecuteResult.success(output);
|
||||||
|
|||||||
@ -37,6 +37,7 @@ public class FlowLoopNodeHandler implements FlowNodeTypeHandler {
|
|||||||
String nodeId = message.getNodeId();
|
String nodeId = message.getNodeId();
|
||||||
FlowNodeWrapper nodeWrapper = graph.getNode(nodeId);
|
FlowNodeWrapper nodeWrapper = graph.getNode(nodeId);
|
||||||
int loopCount = inputParams.getIntValue("loopNum"); // 获取循环次数
|
int loopCount = inputParams.getIntValue("loopNum"); // 获取循环次数
|
||||||
|
Object loopArray = inputParams.get("loopArray"); // 获取循环次数
|
||||||
|
|
||||||
// 如果没有设置循环次数或者循环次数为 0,则跳过处理
|
// 如果没有设置循环次数或者循环次数为 0,则跳过处理
|
||||||
if (loopCount <= 0) {
|
if (loopCount <= 0) {
|
||||||
@ -65,8 +66,9 @@ public class FlowLoopNodeHandler implements FlowNodeTypeHandler {
|
|||||||
// 克隆 message,并明确设置 loopNum
|
// 克隆 message,并明确设置 loopNum
|
||||||
TaskNodeExecuteMessage subMessage = new TaskNodeExecuteMessage();
|
TaskNodeExecuteMessage subMessage = new TaskNodeExecuteMessage();
|
||||||
BeanUtil.copyProperties(message, subMessage);
|
BeanUtil.copyProperties(message, subMessage);
|
||||||
// subMessage.setLoopNum(i); // 当前 loop 第 i 次
|
subMessage.setLoopNum(i); // 当前 loop 第 i 次
|
||||||
subMessage.setIterations(newIterations); // 完整路径
|
subMessage.setIterations(newIterations); // 完整路径
|
||||||
|
subMessage.setLoopArray(loopArray);
|
||||||
|
|
||||||
flowItemExecutor.executeSubGraph(subGraph, subMessage, latch::countDown, newIterations);
|
flowItemExecutor.executeSubGraph(subGraph, subMessage, latch::countDown, newIterations);
|
||||||
|
|
||||||
|
|||||||
@ -73,7 +73,7 @@ public class FlowItemExecutor {
|
|||||||
List<FlowParamDef> startInputDefs = subGraph.getStartInputParamsDefs();
|
List<FlowParamDef> startInputDefs = subGraph.getStartInputParamsDefs();
|
||||||
|
|
||||||
// 子图的 start 节点
|
// 子图的 start 节点
|
||||||
FlowNodeWrapper subStartNode = subGraph.findSubStartNodeId();
|
FlowNodeWrapper subStartNode = subGraph.getInDegreeZeroNode();
|
||||||
String subStartNodeId = subStartNode.getNodeId();
|
String subStartNodeId = subStartNode.getNodeId();
|
||||||
|
|
||||||
NodeExecutor executor = node -> executeNode(
|
NodeExecutor executor = node -> executeNode(
|
||||||
@ -84,7 +84,7 @@ public class FlowItemExecutor {
|
|||||||
// 执行子图起始节点
|
// 执行子图起始节点
|
||||||
executeNode(subGraph, subStartNode, subStartNodeId, startInputDefs,
|
executeNode(subGraph, subStartNode, subStartNodeId, startInputDefs,
|
||||||
rootMessage, onFinished, iterations);
|
rootMessage, onFinished, iterations);
|
||||||
|
log.info("当前执行-子图");
|
||||||
// 启动子图调度
|
// 启动子图调度
|
||||||
flowTaskScheduler.subStart(subGraph, taskInstHolder.getContext(rootMessage.getInstId()),
|
flowTaskScheduler.subStart(subGraph, taskInstHolder.getContext(rootMessage.getInstId()),
|
||||||
executor, onFinished, iterations);
|
executor, onFinished, iterations);
|
||||||
@ -131,7 +131,7 @@ public class FlowItemExecutor {
|
|||||||
}
|
}
|
||||||
|
|
||||||
// END 节点 or 子图结束
|
// 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 {
|
try {
|
||||||
log.info("{} 节点执行,释放线程", node.getNodeType());
|
log.info("{} 节点执行,释放线程", node.getNodeType());
|
||||||
onFinished.run();
|
onFinished.run();
|
||||||
@ -172,6 +172,8 @@ public class FlowItemExecutor {
|
|||||||
// loopNum 不在这里计算,而是由 FlowLoopNodeHandler / resume 显式写入
|
// loopNum 不在这里计算,而是由 FlowLoopNodeHandler / resume 显式写入
|
||||||
// 如果 rootMessage 里已经带了 loopNum,就沿用它
|
// 如果 rootMessage 里已经带了 loopNum,就沿用它
|
||||||
message.setLoopNum(rootMessage.getLoopNum());
|
message.setLoopNum(rootMessage.getLoopNum());
|
||||||
|
|
||||||
|
message.setLoopArray(rootMessage.getLoopArray());
|
||||||
return message;
|
return message;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
@ -58,7 +58,7 @@ public class FlowTaskScheduler {
|
|||||||
Runnable onFinished,
|
Runnable onFinished,
|
||||||
List<Integer> iterations) {
|
List<Integer> iterations) {
|
||||||
String instId = context.getInstId();
|
String instId = context.getInstId();
|
||||||
FlowNodeWrapper startNode = graph.findSubStartNodeId();
|
FlowNodeWrapper startNode = graph.getInDegreeZeroNode();
|
||||||
List<String> nextNodes = graph.getNextNodes(startNode.getNodeId());
|
List<String> nextNodes = graph.getNextNodes(startNode.getNodeId());
|
||||||
|
|
||||||
if (nextNodes.isEmpty()) {
|
if (nextNodes.isEmpty()) {
|
||||||
|
|||||||
@ -84,11 +84,12 @@ public class FlowNodeParamPreparer {
|
|||||||
// 4 LOOP 节点:动态计算 loopCount
|
// 4 LOOP 节点:动态计算 loopCount
|
||||||
if (node.getNodeType() == NodeTypeEnum.LOOP) {
|
if (node.getNodeType() == NodeTypeEnum.LOOP) {
|
||||||
Object loopNumVal = input.get("loopNum");
|
Object loopNumVal = input.get("loopNum");
|
||||||
|
input.put("loopArray", loopNumVal);
|
||||||
int loopCount = 0;
|
int loopCount = 0;
|
||||||
|
|
||||||
// 根据迭代路径找到当前层的集合对象
|
// 根据迭代路径找到当前层的集合对象
|
||||||
Object target = resolveLoopTarget(loopNumVal, rootMessage.getIterations(), 0);
|
// Object target = resolveLoopTarget(loopNumVal, rootMessage.getIterations(), 0);
|
||||||
|
Object target = loopNumVal;
|
||||||
if (target instanceof JSONArray) {
|
if (target instanceof JSONArray) {
|
||||||
loopCount = ((JSONArray) target).size();
|
loopCount = ((JSONArray) target).size();
|
||||||
} else if (target instanceof Collection) {
|
} else if (target instanceof Collection) {
|
||||||
|
|||||||
@ -22,5 +22,6 @@ public class TaskNodeExecuteContext {
|
|||||||
private TaskNodeExecuteMessage rootMessage;
|
private TaskNodeExecuteMessage rootMessage;
|
||||||
private Runnable onFinished;
|
private Runnable onFinished;
|
||||||
private int loopIteration;
|
private int loopIteration;
|
||||||
|
private Object loopArray;
|
||||||
private List<Integer> iterations = new ArrayList<>();
|
private List<Integer> iterations = new ArrayList<>();
|
||||||
}
|
}
|
||||||
|
|||||||
@ -55,6 +55,8 @@ public class TaskNodeExecuteMessage {
|
|||||||
*/
|
*/
|
||||||
private List<Integer> iterations = new ArrayList<>();
|
private List<Integer> iterations = new ArrayList<>();
|
||||||
|
|
||||||
|
private Object loopArray;
|
||||||
|
|
||||||
/**
|
/**
|
||||||
* 节点名称
|
* 节点名称
|
||||||
*/
|
*/
|
||||||
|
|||||||
Loading…
Reference in New Issue
Block a user