feat(workflow): 实现工作流版本管理功能
- 将检测项流程数据存储分离为草稿和发布版本两个概念 - 添加 TeDetectionItemVersion 实体用于存储不可变的工作流发布版本 - 实现工作流草稿保存、版本发布、版本回退等功能 - 在任务执行时验证检测项是否已发布,防止使用未发布流程 - 添加 TeTaskInstFlowVersion 表记录任务实例关联的具体版本清单 - 重构 TeFlowController 控制器接口,支持草稿保存和版本管理操作 - 更新数据库表结构,添加版本控制相关字段和约束
This commit is contained in:
parent
05f695c93c
commit
6e01f6888b
@ -11,13 +11,16 @@ import com.cmvr.test.flow.context.TaskContextManager;
|
||||
import com.cmvr.test.flow.control.FlowControlService;
|
||||
import com.cmvr.test.flow.runtime.engine.FlowTaskRuntimeService;
|
||||
import com.cmvr.test.flow.runtime.operator.edge.ti.TiTouchOperateService;
|
||||
import com.cmvr.test.model.domain.TeDetectionItem;
|
||||
import com.cmvr.test.model.vo.FlowActionRequestVO;
|
||||
import com.cmvr.test.model.vo.TeDetectItemDeployFlowVO;
|
||||
import com.cmvr.test.model.vo.TeFlowPublishVO;
|
||||
import com.cmvr.test.model.vo.TeFlowVersionActionVO;
|
||||
import com.cmvr.test.model.vo.TeFlowViewVO;
|
||||
import com.cmvr.test.model.vo.TeTaskExecuteNormalVO;
|
||||
import com.cmvr.test.model.vo.TeTaskExecuteTrailVO;
|
||||
import com.cmvr.test.service.FlowActionExecutorService;
|
||||
import com.cmvr.test.service.ITeDetectionItemService;
|
||||
import com.cmvr.test.service.ITeDetectionItemVersionService;
|
||||
import com.cmvr.test.service.ITeNodeInstService;
|
||||
import io.swagger.annotations.Api;
|
||||
import io.swagger.annotations.ApiOperation;
|
||||
@ -29,6 +32,7 @@ import org.springframework.web.bind.annotation.RequestBody;
|
||||
import org.springframework.web.bind.annotation.RequestMapping;
|
||||
import org.springframework.web.bind.annotation.RequestParam;
|
||||
import org.springframework.web.bind.annotation.RestController;
|
||||
import org.springframework.web.bind.annotation.PutMapping;
|
||||
|
||||
import jakarta.annotation.PreDestroy;
|
||||
import jakarta.validation.Valid;
|
||||
@ -62,6 +66,7 @@ public class TeFlowController extends BaseController {
|
||||
private static final String AI_AGENT_API_KEY = "d9etum54shheenol3apg";
|
||||
|
||||
private final ITeDetectionItemService teDetectionItemService;
|
||||
private final ITeDetectionItemVersionService detectionItemVersionService;
|
||||
private final ITeNodeInstService nodeInstService;
|
||||
private final FlowTaskRuntimeService flowTaskRuntimeService;
|
||||
private final TaskContextManager taskContextManager;
|
||||
@ -87,8 +92,42 @@ public class TeFlowController extends BaseController {
|
||||
|
||||
@ApiOperation("流程发布")
|
||||
@PostMapping("/publish")
|
||||
public AjaxResult publish(@RequestBody TeDetectionItem detectionItem) {
|
||||
return toAjax(teDetectionItemService.publish(detectionItem.getId()));
|
||||
public AjaxResult publish(@Valid @RequestBody TeFlowPublishVO request) {
|
||||
return success(detectionItemVersionService.publish(request, getUsername()));
|
||||
}
|
||||
|
||||
@ApiOperation("保存工作流草稿")
|
||||
@PutMapping("/draft")
|
||||
public AjaxResult saveDraft(@Valid @RequestBody TeDetectItemDeployFlowVO request) {
|
||||
return success(teDetectionItemService.saveDraft(request, getUsername()));
|
||||
}
|
||||
|
||||
@ApiOperation("查询工作流发布记录")
|
||||
@GetMapping("/versions/{itemId}")
|
||||
public AjaxResult versions(@PathVariable String itemId) {
|
||||
return success(detectionItemVersionService.listVersions(itemId));
|
||||
}
|
||||
|
||||
@ApiOperation("查看工作流发布版本")
|
||||
@GetMapping("/version/{versionId}")
|
||||
public AjaxResult version(@PathVariable String versionId) {
|
||||
return success(detectionItemVersionService.getVersion(versionId));
|
||||
}
|
||||
|
||||
@ApiOperation("将发布版本恢复为草稿")
|
||||
@PostMapping("/version/{versionId}/restore-draft")
|
||||
public AjaxResult restoreDraft(@PathVariable String versionId,
|
||||
@RequestBody(required = false) TeFlowVersionActionVO request) {
|
||||
Integer revision = request == null ? null : request.getDraftRevision();
|
||||
return success(detectionItemVersionService.restoreDraft(versionId, revision, getUsername()));
|
||||
}
|
||||
|
||||
@ApiOperation("回退到指定发布版本")
|
||||
@PostMapping("/version/{versionId}/rollback")
|
||||
public AjaxResult rollback(@PathVariable String versionId,
|
||||
@RequestBody(required = false) TeFlowVersionActionVO request) {
|
||||
String note = request == null ? null : request.getPublishNote();
|
||||
return success(detectionItemVersionService.rollback(versionId, note, getUsername()));
|
||||
}
|
||||
|
||||
@ApiOperation("演示任务执行入口")
|
||||
|
||||
@ -14,13 +14,13 @@ import com.cmvr.test.flow.context.TaskContext;
|
||||
import com.cmvr.test.flow.context.TaskInstHolder;
|
||||
import com.cmvr.test.model.domain.TeTaskInst;
|
||||
import com.cmvr.test.model.dto.TeQueryTaskDetailDTO;
|
||||
import com.cmvr.test.model.vo.TeDetectItemDeployFlowVO;
|
||||
import com.cmvr.test.model.domain.TeTaskInstFlowVersion;
|
||||
import com.cmvr.test.model.vo.TeTaskExecuteNormalVO;
|
||||
import com.cmvr.test.model.vo.TeTaskExecuteProjectVO;
|
||||
import com.cmvr.test.model.vo.TeTaskExecuteTrailVO;
|
||||
import com.cmvr.test.model.vo.ResourceRuntimeAssignmentVO;
|
||||
import com.cmvr.test.model.vo.RobotRuntimeMemberVO;
|
||||
import com.cmvr.test.service.ITeDetectionItemService;
|
||||
import com.cmvr.test.service.ITeTaskInstFlowVersionService;
|
||||
import com.cmvr.test.service.ITeTaskInstService;
|
||||
import com.cmvr.test.service.ITeTaskOrchestrationService;
|
||||
import lombok.RequiredArgsConstructor;
|
||||
@ -48,7 +48,7 @@ public class FlowTaskRuntimeEntry implements FlowTaskRuntimeService {
|
||||
|
||||
private final ITeTaskOrchestrationService taskOrchestrationService;
|
||||
private final ITeTaskInstService taskInstService;
|
||||
private final ITeDetectionItemService detectionItemService;
|
||||
private final ITeTaskInstFlowVersionService taskInstFlowVersionService;
|
||||
private final TaskInstHolder taskInstHolder;
|
||||
private final FlowTaskAsyncDispatcher flowTaskAsyncDispatcher;
|
||||
private final ObjectProvider<RobotExecutionTargetResolver> targetResolverProvider;
|
||||
@ -127,11 +127,6 @@ public class FlowTaskRuntimeEntry implements FlowTaskRuntimeService {
|
||||
String itemId = taskExecuteTrailVO.getItemId();
|
||||
runtimeDefinitionValidator.validate(taskExecuteTrailVO.getFlowData(), resources.assignments());
|
||||
|
||||
TeDetectItemDeployFlowVO deployFlowVO = new TeDetectItemDeployFlowVO();
|
||||
deployFlowVO.setId(itemId);
|
||||
deployFlowVO.setFlowData(taskExecuteTrailVO.getFlowData());
|
||||
detectionItemService.deployFlow(deployFlowVO);
|
||||
|
||||
String taskId = taskOrchestrationService.buildTrialTask(itemId);
|
||||
|
||||
// 创建并保存任务实例
|
||||
@ -196,6 +191,12 @@ public class FlowTaskRuntimeEntry implements FlowTaskRuntimeService {
|
||||
List<TeQueryTaskDetailDTO> details = new ArrayList<>(taskOrchestrationService.queryTaskDetail(taskId));
|
||||
details.sort(Comparator.comparing(TeQueryTaskDetailDTO::getOrderNum,
|
||||
Comparator.nullsLast(Integer::compareTo)));
|
||||
TeQueryTaskDetailDTO unpublished = details.stream()
|
||||
.filter(detail -> StrUtil.isBlank(detail.getVersionId()) || StrUtil.isBlank(detail.getFlowData()))
|
||||
.findFirst().orElse(null);
|
||||
if (unpublished != null) {
|
||||
throw new GlobalException("任务中的检测项尚未发布:" + unpublished.getDetectItemId());
|
||||
}
|
||||
Map<String, Integer> validationOccurrences = new HashMap<>();
|
||||
details.forEach(detail -> {
|
||||
int occurrence = validationOccurrences.merge(detail.getDetectItemId(), 1, Integer::sum);
|
||||
@ -203,13 +204,10 @@ public class FlowTaskRuntimeEntry implements FlowTaskRuntimeService {
|
||||
FlowRuntimeAssignmentScope.forItem(
|
||||
resources.assignments(), detail.getDetectItemId(), occurrence));
|
||||
});
|
||||
int count = (int) details.stream().filter(dto -> ObjUtil.equals(dto.getIsDeploy(), "1")).count();
|
||||
if (count > 0) {
|
||||
throw new GlobalException("当前任务存在未发布的检测项!");
|
||||
}
|
||||
// 创建任务实例
|
||||
TeTaskInst instance = createAndSaveTaskInstance(taskId, target, runMode, originalVO);
|
||||
String instId = instance.getId();
|
||||
saveFlowVersionManifest(instId, details);
|
||||
|
||||
try {
|
||||
// 注册上下文
|
||||
@ -229,6 +227,25 @@ public class FlowTaskRuntimeEntry implements FlowTaskRuntimeService {
|
||||
}
|
||||
}
|
||||
|
||||
private void saveFlowVersionManifest(String instId, List<TeQueryTaskDetailDTO> details) {
|
||||
Map<String, Integer> occurrences = new HashMap<>();
|
||||
List<TeTaskInstFlowVersion> manifest = new ArrayList<>();
|
||||
for (int i = 0; i < details.size(); i++) {
|
||||
TeQueryTaskDetailDTO detail = details.get(i);
|
||||
TeTaskInstFlowVersion entry = new TeTaskInstFlowVersion();
|
||||
entry.setInstId(instId);
|
||||
entry.setItemId(detail.getDetectItemId());
|
||||
entry.setItemOrder(detail.getOrderNum() == null ? i + 1 : detail.getOrderNum());
|
||||
entry.setItemOccurrence(occurrences.merge(detail.getDetectItemId(), 1, Integer::sum));
|
||||
entry.setVersionId(detail.getVersionId());
|
||||
entry.setVersionNo(detail.getVersionNo());
|
||||
manifest.add(entry);
|
||||
}
|
||||
if (!manifest.isEmpty() && !taskInstFlowVersionService.saveBatch(manifest)) {
|
||||
throw new GlobalException("保存任务工作流版本清单失败");
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* 创建并保存任务实例
|
||||
*/
|
||||
|
||||
@ -0,0 +1,7 @@
|
||||
package com.cmvr.test.mapper;
|
||||
|
||||
import com.baomidou.mybatisplus.core.mapper.BaseMapper;
|
||||
import com.cmvr.test.model.domain.TeDetectionItemVersion;
|
||||
|
||||
public interface TeDetectionItemVersionMapper extends BaseMapper<TeDetectionItemVersion> {
|
||||
}
|
||||
@ -0,0 +1,7 @@
|
||||
package com.cmvr.test.mapper;
|
||||
|
||||
import com.baomidou.mybatisplus.core.mapper.BaseMapper;
|
||||
import com.cmvr.test.model.domain.TeTaskInstFlowVersion;
|
||||
|
||||
public interface TeTaskInstFlowVersionMapper extends BaseMapper<TeTaskInstFlowVersion> {
|
||||
}
|
||||
@ -36,6 +36,21 @@ public class TeDetectionItem extends BaseEntity {
|
||||
@ApiModelProperty("检测项流程数据")
|
||||
private String flowData;
|
||||
|
||||
@ApiModelProperty("当前发布版本ID")
|
||||
private String publishedVersionId;
|
||||
|
||||
@ApiModelProperty("当前发布版本号")
|
||||
private Integer publishedVersionNo;
|
||||
|
||||
@ApiModelProperty("草稿修订号")
|
||||
private Integer draftRevision;
|
||||
|
||||
@ApiModelProperty("草稿内容摘要")
|
||||
private String draftHash;
|
||||
|
||||
@ApiModelProperty("最近发布的草稿内容摘要")
|
||||
private String publishedDraftHash;
|
||||
|
||||
@Excel(name = "状态")
|
||||
@ApiModelProperty(value = "状态 (0正常 1停用)")
|
||||
@TableField(fill = FieldFill.INSERT)
|
||||
|
||||
@ -0,0 +1,42 @@
|
||||
package com.cmvr.test.model.domain;
|
||||
|
||||
import com.baomidou.mybatisplus.annotation.IdType;
|
||||
import com.baomidou.mybatisplus.annotation.TableId;
|
||||
import com.cmvr.common.core.domain.BaseEntity;
|
||||
import io.swagger.annotations.ApiModel;
|
||||
import io.swagger.annotations.ApiModelProperty;
|
||||
import lombok.Data;
|
||||
import lombok.EqualsAndHashCode;
|
||||
|
||||
@Data
|
||||
@EqualsAndHashCode(callSuper = true)
|
||||
@ApiModel("工作流发布版本")
|
||||
public class TeDetectionItemVersion extends BaseEntity {
|
||||
|
||||
@TableId(value = "id", type = IdType.ASSIGN_UUID)
|
||||
private String id;
|
||||
|
||||
@ApiModelProperty("检测项ID")
|
||||
private String detectionItemId;
|
||||
|
||||
@ApiModelProperty("发布版本号")
|
||||
private Integer versionNo;
|
||||
|
||||
@ApiModelProperty("不可变流程定义")
|
||||
private String flowData;
|
||||
|
||||
@ApiModelProperty("不可变运行参数定义")
|
||||
private String config;
|
||||
|
||||
@ApiModelProperty("流程JSON结构版本")
|
||||
private Integer schemaVersion;
|
||||
|
||||
@ApiModelProperty("发布内容摘要")
|
||||
private String contentHash;
|
||||
|
||||
@ApiModelProperty("回退来源版本ID")
|
||||
private String sourceVersionId;
|
||||
|
||||
@ApiModelProperty("发布说明")
|
||||
private String publishNote;
|
||||
}
|
||||
@ -0,0 +1,21 @@
|
||||
package com.cmvr.test.model.domain;
|
||||
|
||||
import com.baomidou.mybatisplus.annotation.IdType;
|
||||
import com.baomidou.mybatisplus.annotation.TableId;
|
||||
import com.cmvr.common.core.domain.BaseEntity;
|
||||
import lombok.Data;
|
||||
import lombok.EqualsAndHashCode;
|
||||
|
||||
@Data
|
||||
@EqualsAndHashCode(callSuper = true)
|
||||
public class TeTaskInstFlowVersion extends BaseEntity {
|
||||
|
||||
@TableId(value = "id", type = IdType.ASSIGN_UUID)
|
||||
private String id;
|
||||
private String instId;
|
||||
private String itemId;
|
||||
private Integer itemOrder;
|
||||
private Integer itemOccurrence;
|
||||
private String versionId;
|
||||
private Integer versionNo;
|
||||
}
|
||||
@ -8,6 +8,8 @@ public class TeQueryTaskDetailDTO {
|
||||
private String detectItemId;
|
||||
private String flowData;
|
||||
private String config;
|
||||
private String versionId;
|
||||
private Integer versionNo;
|
||||
private Integer orderNum;
|
||||
private String isDeploy;
|
||||
private String sceneType;
|
||||
|
||||
@ -21,4 +21,7 @@ public class TeDetectItemDeployFlowVO {
|
||||
|
||||
@ApiModelProperty(value = "流程运行参数")
|
||||
private JSONObject config;
|
||||
|
||||
@ApiModelProperty("客户端加载草稿时的修订号")
|
||||
private Integer draftRevision;
|
||||
}
|
||||
|
||||
@ -0,0 +1,20 @@
|
||||
package com.cmvr.test.model.vo;
|
||||
|
||||
import io.swagger.annotations.ApiModel;
|
||||
import io.swagger.annotations.ApiModelProperty;
|
||||
import jakarta.validation.constraints.NotBlank;
|
||||
import lombok.Data;
|
||||
|
||||
@Data
|
||||
@ApiModel("工作流发布请求")
|
||||
public class TeFlowPublishVO {
|
||||
|
||||
@NotBlank(message = "检测项ID不能为空")
|
||||
private String id;
|
||||
|
||||
@ApiModelProperty("发布所基于的草稿修订号")
|
||||
private Integer draftRevision;
|
||||
|
||||
@ApiModelProperty("发布说明")
|
||||
private String publishNote;
|
||||
}
|
||||
@ -0,0 +1,12 @@
|
||||
package com.cmvr.test.model.vo;
|
||||
|
||||
import io.swagger.annotations.ApiModelProperty;
|
||||
import lombok.Data;
|
||||
|
||||
@Data
|
||||
public class TeFlowVersionActionVO {
|
||||
@ApiModelProperty("客户端加载草稿时的修订号")
|
||||
private Integer draftRevision;
|
||||
@ApiModelProperty("回退说明")
|
||||
private String publishNote;
|
||||
}
|
||||
@ -16,6 +16,12 @@ public class TeQueryTaskOrchestraItemVO {
|
||||
@ApiModelProperty("检测项参数信息")
|
||||
private String config;
|
||||
|
||||
@ApiModelProperty("当前发布版本ID")
|
||||
private String versionId;
|
||||
|
||||
@ApiModelProperty("当前发布版本号")
|
||||
private Integer versionNo;
|
||||
|
||||
@ApiModelProperty("执行顺序号")
|
||||
private Integer orderNum;
|
||||
}
|
||||
|
||||
@ -58,6 +58,9 @@ public interface ITeDetectionItemService extends IService<TeDetectionItem> {
|
||||
*/
|
||||
public int deployFlow(TeDetectItemDeployFlowVO teDeployFlowVo);
|
||||
|
||||
/** Save the editable workflow draft without changing the published release. */
|
||||
TeDetectionItem saveDraft(TeDetectItemDeployFlowVO request, String operator);
|
||||
|
||||
/**
|
||||
* 发布流程
|
||||
*/
|
||||
|
||||
@ -0,0 +1,16 @@
|
||||
package com.cmvr.test.service;
|
||||
|
||||
import com.baomidou.mybatisplus.spring.service.IService;
|
||||
import com.cmvr.test.model.domain.TeDetectionItem;
|
||||
import com.cmvr.test.model.domain.TeDetectionItemVersion;
|
||||
import com.cmvr.test.model.vo.TeFlowPublishVO;
|
||||
|
||||
import java.util.List;
|
||||
|
||||
public interface ITeDetectionItemVersionService extends IService<TeDetectionItemVersion> {
|
||||
TeDetectionItemVersion publish(TeFlowPublishVO request, String operator);
|
||||
List<TeDetectionItemVersion> listVersions(String itemId);
|
||||
TeDetectionItemVersion getVersion(String versionId);
|
||||
TeDetectionItem restoreDraft(String versionId, Integer expectedRevision, String operator);
|
||||
TeDetectionItemVersion rollback(String versionId, String publishNote, String operator);
|
||||
}
|
||||
@ -0,0 +1,7 @@
|
||||
package com.cmvr.test.service;
|
||||
|
||||
import com.baomidou.mybatisplus.spring.service.IService;
|
||||
import com.cmvr.test.model.domain.TeTaskInstFlowVersion;
|
||||
|
||||
public interface ITeTaskInstFlowVersionService extends IService<TeTaskInstFlowVersion> {
|
||||
}
|
||||
@ -1,13 +1,15 @@
|
||||
package com.cmvr.test.service.impl;
|
||||
|
||||
import cn.hutool.core.util.StrUtil;
|
||||
import cn.hutool.crypto.digest.DigestUtil;
|
||||
import com.alibaba.fastjson2.JSON;
|
||||
import com.alibaba.fastjson2.JSONArray;
|
||||
import com.alibaba.fastjson2.JSONObject;
|
||||
import com.baomidou.mybatisplus.core.conditions.query.LambdaQueryWrapper;
|
||||
import com.baomidou.mybatisplus.core.conditions.query.QueryWrapper;
|
||||
import com.baomidou.mybatisplus.core.conditions.update.LambdaUpdateWrapper;
|
||||
import com.baomidou.mybatisplus.spring.service.impl.ServiceImpl;
|
||||
import com.cmvr.common.core.domain.entity.SysGroup;
|
||||
import com.cmvr.common.exception.GlobalException;
|
||||
import com.cmvr.common.utils.EntityCopyUtils;
|
||||
import com.cmvr.test.enums.NodeTypeEnum;
|
||||
import com.cmvr.test.mapper.TeDetectionItemMapper;
|
||||
@ -15,6 +17,7 @@ import com.cmvr.test.model.domain.TeDetectionItem;
|
||||
import com.cmvr.test.model.vo.TeDetectItemDeployFlowVO;
|
||||
import com.cmvr.test.model.vo.TeDetectionItemVO;
|
||||
import com.cmvr.test.service.ITeDetectionItemService;
|
||||
import com.cmvr.test.service.ITeDetectionItemVersionService;
|
||||
import com.github.yulichang.wrapper.MPJLambdaWrapper;
|
||||
import lombok.RequiredArgsConstructor;
|
||||
import org.springframework.stereotype.Service;
|
||||
@ -33,6 +36,8 @@ import java.util.List;
|
||||
@RequiredArgsConstructor
|
||||
public class TeDetectionItemServiceImpl extends ServiceImpl<TeDetectionItemMapper, TeDetectionItem> implements ITeDetectionItemService {
|
||||
|
||||
private final ITeDetectionItemVersionService versionService;
|
||||
|
||||
/**
|
||||
* 查询检测项配置
|
||||
*
|
||||
@ -81,6 +86,7 @@ public class TeDetectionItemServiceImpl extends ServiceImpl<TeDetectionItemMappe
|
||||
// 新增基础配置 此时状态和部署 皆为不可用
|
||||
teDetectionItem.setStatus("1");
|
||||
teDetectionItem.setIsDeploy("1");
|
||||
teDetectionItem.setDraftRevision(0);
|
||||
return this.baseMapper.insert(teDetectionItem);
|
||||
}
|
||||
|
||||
@ -97,12 +103,14 @@ public class TeDetectionItemServiceImpl extends ServiceImpl<TeDetectionItemMappe
|
||||
.detectDesc(source.getDetectDesc())
|
||||
.flowData(source.getFlowData())
|
||||
.status(source.getStatus())
|
||||
.isDeploy(source.getIsDeploy())
|
||||
.isDeploy("1")
|
||||
.detectVersion(source.getDetectVersion())
|
||||
.groupId(source.getGroupId())
|
||||
.config(source.getConfig())
|
||||
.sceneCode(source.getSceneCode())
|
||||
.build();
|
||||
copy.setDraftRevision(0);
|
||||
copy.setDraftHash(contentHash(copy.getFlowData(), copy.getConfig()));
|
||||
copy.setRemark(source.getRemark());
|
||||
this.baseMapper.insert(copy);
|
||||
return copy.getId();
|
||||
@ -125,21 +133,41 @@ public class TeDetectionItemServiceImpl extends ServiceImpl<TeDetectionItemMappe
|
||||
@Override
|
||||
@Transactional
|
||||
public int deployFlow(TeDetectItemDeployFlowVO teDeployFlowVo) {
|
||||
// 校验记录存在
|
||||
boolean exists = this.baseMapper
|
||||
.exists(new QueryWrapper<TeDetectionItem>()
|
||||
.eq("id", teDeployFlowVo.getId()));
|
||||
saveDraft(teDeployFlowVo, null);
|
||||
return 1;
|
||||
}
|
||||
|
||||
if (!exists) {
|
||||
throw new IllegalArgumentException("ID为" + teDeployFlowVo.getId() + "的检测项不存在!");
|
||||
@Override
|
||||
@Transactional(rollbackFor = Exception.class)
|
||||
public TeDetectionItem saveDraft(TeDetectItemDeployFlowVO request, String operator) {
|
||||
TeDetectionItem current = this.baseMapper.selectById(request.getId());
|
||||
if (current == null) {
|
||||
throw new IllegalArgumentException("ID为" + request.getId() + "的检测项不存在!");
|
||||
}
|
||||
int currentRevision = current.getDraftRevision() == null ? 0 : current.getDraftRevision();
|
||||
if (request.getDraftRevision() != null && request.getDraftRevision() != currentRevision) {
|
||||
throw new GlobalException("草稿已被其他窗口修改,请刷新后再保存");
|
||||
}
|
||||
|
||||
// 保存流程 JSON
|
||||
String config = extractStartConfig(request.getFlowData());
|
||||
TeDetectionItem updateData = new TeDetectionItem();
|
||||
updateData.setId(teDeployFlowVo.getId());
|
||||
updateData.setFlowData(teDeployFlowVo.getFlowData());
|
||||
//获取params
|
||||
JSONObject flow = JSON.parseObject(teDeployFlowVo.getFlowData());
|
||||
updateData.setFlowData(request.getFlowData());
|
||||
updateData.setConfig(config);
|
||||
updateData.setDraftRevision(currentRevision + 1);
|
||||
updateData.setDraftHash(contentHash(request.getFlowData(), config));
|
||||
updateData.setUpdateBy(operator);
|
||||
int rows = this.baseMapper.update(updateData,
|
||||
new LambdaUpdateWrapper<TeDetectionItem>()
|
||||
.eq(TeDetectionItem::getId, request.getId())
|
||||
.eq(TeDetectionItem::getDraftRevision, currentRevision));
|
||||
if (rows != 1) {
|
||||
throw new GlobalException("草稿已被其他窗口修改,请刷新后再保存");
|
||||
}
|
||||
return this.baseMapper.selectById(request.getId());
|
||||
}
|
||||
|
||||
private String extractStartConfig(String flowData) {
|
||||
JSONObject flow = JSON.parseObject(flowData);
|
||||
JSONArray nodeArray = flow.getJSONArray("nodes");
|
||||
if (nodeArray == null) {
|
||||
throw new IllegalArgumentException("流程图中未找到任何节点");
|
||||
@ -170,20 +198,22 @@ public class TeDetectionItemServiceImpl extends ServiceImpl<TeDetectionItemMappe
|
||||
}
|
||||
}
|
||||
|
||||
updateData.setConfig(inputParams.toJSONString());
|
||||
break;
|
||||
return inputParams.toJSONString();
|
||||
}
|
||||
return "[]";
|
||||
}
|
||||
|
||||
return this.baseMapper.updateById(updateData);
|
||||
private String contentHash(String flowData, String config) {
|
||||
return StrUtil.isBlank(flowData) ? null
|
||||
: DigestUtil.sha256Hex(flowData + "\n" + StrUtil.nullToEmpty(config));
|
||||
}
|
||||
|
||||
@Override
|
||||
public int publish(String flowId) {
|
||||
TeDetectionItem updateData = new TeDetectionItem();
|
||||
updateData.setId(flowId);
|
||||
updateData.setIsDeploy("0");
|
||||
updateData.setStatus("0");
|
||||
return this.baseMapper.updateById(updateData);
|
||||
com.cmvr.test.model.vo.TeFlowPublishVO request = new com.cmvr.test.model.vo.TeFlowPublishVO();
|
||||
request.setId(flowId);
|
||||
versionService.publish(request, null);
|
||||
return 1;
|
||||
}
|
||||
|
||||
|
||||
|
||||
@ -0,0 +1,179 @@
|
||||
package com.cmvr.test.service.impl;
|
||||
|
||||
import cn.hutool.core.util.StrUtil;
|
||||
import cn.hutool.crypto.digest.DigestUtil;
|
||||
import com.alibaba.fastjson2.JSON;
|
||||
import com.alibaba.fastjson2.JSONObject;
|
||||
import com.baomidou.mybatisplus.core.conditions.query.LambdaQueryWrapper;
|
||||
import com.baomidou.mybatisplus.core.conditions.update.LambdaUpdateWrapper;
|
||||
import com.baomidou.mybatisplus.spring.service.impl.ServiceImpl;
|
||||
import com.cmvr.common.exception.GlobalException;
|
||||
import com.cmvr.test.mapper.TeDetectionItemMapper;
|
||||
import com.cmvr.test.mapper.TeDetectionItemVersionMapper;
|
||||
import com.cmvr.test.model.domain.TeDetectionItem;
|
||||
import com.cmvr.test.model.domain.TeDetectionItemVersion;
|
||||
import com.cmvr.test.model.vo.TeFlowPublishVO;
|
||||
import com.cmvr.test.service.ITeDetectionItemVersionService;
|
||||
import lombok.RequiredArgsConstructor;
|
||||
import org.springframework.stereotype.Service;
|
||||
import org.springframework.transaction.annotation.Transactional;
|
||||
|
||||
import java.util.List;
|
||||
|
||||
@Service
|
||||
@RequiredArgsConstructor
|
||||
public class TeDetectionItemVersionServiceImpl
|
||||
extends ServiceImpl<TeDetectionItemVersionMapper, TeDetectionItemVersion>
|
||||
implements ITeDetectionItemVersionService {
|
||||
|
||||
private final TeDetectionItemMapper detectionItemMapper;
|
||||
|
||||
@Override
|
||||
@Transactional(rollbackFor = Exception.class)
|
||||
public TeDetectionItemVersion publish(TeFlowPublishVO request, String operator) {
|
||||
TeDetectionItem item = requireItem(request.getId(), true);
|
||||
assertRevision(item, request.getDraftRevision());
|
||||
if (StrUtil.isBlank(item.getFlowData())) {
|
||||
throw new GlobalException("请先保存工作流草稿");
|
||||
}
|
||||
String hash = contentHash(item.getFlowData(), item.getConfig());
|
||||
TeDetectionItemVersion current = StrUtil.isBlank(item.getPublishedVersionId())
|
||||
? null : baseMapper.selectById(item.getPublishedVersionId());
|
||||
if (current != null && hash.equals(current.getContentHash())) {
|
||||
return current;
|
||||
}
|
||||
|
||||
TeDetectionItemVersion version = createVersion(item, nextVersionNo(item.getId()),
|
||||
item.getFlowData(), item.getConfig(), hash, null, request.getPublishNote(), operator);
|
||||
baseMapper.insert(version);
|
||||
pointToVersion(item.getId(), version, hash, operator);
|
||||
return version;
|
||||
}
|
||||
|
||||
@Override
|
||||
public List<TeDetectionItemVersion> listVersions(String itemId) {
|
||||
return baseMapper.selectList(new LambdaQueryWrapper<TeDetectionItemVersion>()
|
||||
.select(TeDetectionItemVersion::getId, TeDetectionItemVersion::getDetectionItemId,
|
||||
TeDetectionItemVersion::getVersionNo, TeDetectionItemVersion::getSchemaVersion,
|
||||
TeDetectionItemVersion::getContentHash, TeDetectionItemVersion::getSourceVersionId,
|
||||
TeDetectionItemVersion::getPublishNote, TeDetectionItemVersion::getCreateBy,
|
||||
TeDetectionItemVersion::getCreateTime)
|
||||
.eq(TeDetectionItemVersion::getDetectionItemId, itemId)
|
||||
.orderByDesc(TeDetectionItemVersion::getVersionNo));
|
||||
}
|
||||
|
||||
@Override
|
||||
public TeDetectionItemVersion getVersion(String versionId) {
|
||||
TeDetectionItemVersion version = baseMapper.selectById(versionId);
|
||||
if (version == null) throw new GlobalException("发布版本不存在");
|
||||
return version;
|
||||
}
|
||||
|
||||
@Override
|
||||
@Transactional(rollbackFor = Exception.class)
|
||||
public TeDetectionItem restoreDraft(String versionId, Integer expectedRevision, String operator) {
|
||||
TeDetectionItemVersion version = getVersion(versionId);
|
||||
TeDetectionItem item = requireItem(version.getDetectionItemId(), false);
|
||||
assertRevision(item, expectedRevision);
|
||||
int revision = item.getDraftRevision() == null ? 0 : item.getDraftRevision();
|
||||
TeDetectionItem update = new TeDetectionItem();
|
||||
update.setFlowData(version.getFlowData());
|
||||
update.setConfig(version.getConfig());
|
||||
update.setDraftHash(version.getContentHash());
|
||||
update.setDraftRevision(revision + 1);
|
||||
update.setUpdateBy(operator);
|
||||
int rows = detectionItemMapper.update(update, new LambdaUpdateWrapper<TeDetectionItem>()
|
||||
.eq(TeDetectionItem::getId, item.getId())
|
||||
.eq(TeDetectionItem::getDraftRevision, revision));
|
||||
if (rows != 1) throw new GlobalException("草稿已被其他窗口修改,请刷新后重试");
|
||||
return detectionItemMapper.selectById(item.getId());
|
||||
}
|
||||
|
||||
@Override
|
||||
@Transactional(rollbackFor = Exception.class)
|
||||
public TeDetectionItemVersion rollback(String versionId, String publishNote, String operator) {
|
||||
TeDetectionItemVersion source = getVersion(versionId);
|
||||
TeDetectionItem item = requireItem(source.getDetectionItemId(), true);
|
||||
String note = StrUtil.blankToDefault(publishNote, "回退到 V" + source.getVersionNo());
|
||||
TeDetectionItemVersion version = createVersion(item, nextVersionNo(item.getId()),
|
||||
source.getFlowData(), source.getConfig(), source.getContentHash(), source.getId(), note, operator);
|
||||
baseMapper.insert(version);
|
||||
pointToVersion(item.getId(), version, source.getContentHash(), operator);
|
||||
return version;
|
||||
}
|
||||
|
||||
private TeDetectionItem requireItem(String itemId, boolean lock) {
|
||||
LambdaQueryWrapper<TeDetectionItem> query = new LambdaQueryWrapper<TeDetectionItem>()
|
||||
.eq(TeDetectionItem::getId, itemId);
|
||||
if (lock) query.last("for update");
|
||||
TeDetectionItem item = detectionItemMapper.selectOne(query);
|
||||
if (item == null) throw new GlobalException("检测项不存在");
|
||||
return item;
|
||||
}
|
||||
|
||||
private void assertRevision(TeDetectionItem item, Integer expected) {
|
||||
int actual = item.getDraftRevision() == null ? 0 : item.getDraftRevision();
|
||||
if (expected != null && expected != actual) {
|
||||
throw new GlobalException("草稿已被其他窗口修改,请刷新后重试");
|
||||
}
|
||||
}
|
||||
|
||||
private int nextVersionNo(String itemId) {
|
||||
TeDetectionItemVersion latest = baseMapper.selectOne(new LambdaQueryWrapper<TeDetectionItemVersion>()
|
||||
.eq(TeDetectionItemVersion::getDetectionItemId, itemId)
|
||||
.orderByDesc(TeDetectionItemVersion::getVersionNo)
|
||||
.last("limit 1"));
|
||||
return latest == null ? 1 : latest.getVersionNo() + 1;
|
||||
}
|
||||
|
||||
private TeDetectionItemVersion createVersion(TeDetectionItem item, int versionNo,
|
||||
String flowData, String config, String hash,
|
||||
String sourceVersionId, String note, String operator) {
|
||||
TeDetectionItemVersion version = new TeDetectionItemVersion();
|
||||
version.setDetectionItemId(item.getId());
|
||||
version.setVersionNo(versionNo);
|
||||
version.setFlowData(flowData);
|
||||
version.setConfig(config);
|
||||
version.setSchemaVersion(schemaVersion(flowData));
|
||||
version.setContentHash(hash);
|
||||
version.setSourceVersionId(sourceVersionId);
|
||||
version.setPublishNote(note);
|
||||
version.setCreateBy(operator);
|
||||
return version;
|
||||
}
|
||||
|
||||
private int schemaVersion(String flowData) {
|
||||
try {
|
||||
JSONObject flow = JSON.parseObject(flowData);
|
||||
Integer version = flow.getInteger("schemaVersion");
|
||||
if (version != null) return version;
|
||||
if (flow.getJSONArray("nodes") != null) {
|
||||
for (Object value : flow.getJSONArray("nodes")) {
|
||||
JSONObject node = (JSONObject) value;
|
||||
if ("start".equalsIgnoreCase(node.getString("type")) && node.getJSONObject("properties") != null) {
|
||||
return node.getJSONObject("properties").getIntValue("schemaVersion", 1);
|
||||
}
|
||||
}
|
||||
}
|
||||
return 1;
|
||||
} catch (RuntimeException ignored) {
|
||||
return 1;
|
||||
}
|
||||
}
|
||||
|
||||
private String contentHash(String flowData, String config) {
|
||||
return DigestUtil.sha256Hex(flowData + "\n" + StrUtil.nullToEmpty(config));
|
||||
}
|
||||
|
||||
private void pointToVersion(String itemId, TeDetectionItemVersion version, String hash, String operator) {
|
||||
TeDetectionItem update = new TeDetectionItem();
|
||||
update.setPublishedVersionId(version.getId());
|
||||
update.setPublishedVersionNo(version.getVersionNo());
|
||||
update.setPublishedDraftHash(hash);
|
||||
update.setIsDeploy("0");
|
||||
update.setStatus("0");
|
||||
update.setUpdateBy(operator);
|
||||
detectionItemMapper.update(update, new LambdaUpdateWrapper<TeDetectionItem>()
|
||||
.eq(TeDetectionItem::getId, itemId));
|
||||
}
|
||||
}
|
||||
@ -0,0 +1,13 @@
|
||||
package com.cmvr.test.service.impl;
|
||||
|
||||
import com.baomidou.mybatisplus.spring.service.impl.ServiceImpl;
|
||||
import com.cmvr.test.mapper.TeTaskInstFlowVersionMapper;
|
||||
import com.cmvr.test.model.domain.TeTaskInstFlowVersion;
|
||||
import com.cmvr.test.service.ITeTaskInstFlowVersionService;
|
||||
import org.springframework.stereotype.Service;
|
||||
|
||||
@Service
|
||||
public class TeTaskInstFlowVersionServiceImpl
|
||||
extends ServiceImpl<TeTaskInstFlowVersionMapper, TeTaskInstFlowVersion>
|
||||
implements ITeTaskInstFlowVersionService {
|
||||
}
|
||||
@ -8,6 +8,7 @@ import com.cmvr.common.utils.uuid.IdUtils;
|
||||
import com.cmvr.test.enums.RunModeEnum;
|
||||
import com.cmvr.test.mapper.TeTaskOrchestrationMapper;
|
||||
import com.cmvr.test.model.domain.TeDetectionItem;
|
||||
import com.cmvr.test.model.domain.TeDetectionItemVersion;
|
||||
import com.cmvr.test.model.domain.TeTaskConfigInfo;
|
||||
import com.cmvr.test.model.domain.TeTaskOrchestration;
|
||||
import com.cmvr.test.model.dto.TeQueryTaskDetailDTO;
|
||||
@ -43,7 +44,10 @@ public class TeTaskOrchestrationServiceImpl extends ServiceImpl<TeTaskOrchestrat
|
||||
.select(TeTaskOrchestration::getOrderNum)
|
||||
.leftJoin(TeDetectionItem.class, TeDetectionItem::getId, TeTaskOrchestration::getItemId)
|
||||
.selectAs(TeDetectionItem::getDetectName, TeQueryTaskOrchestraItemVO::getItemName)
|
||||
.selectAs(TeDetectionItem::getConfig, TeQueryTaskOrchestraItemVO::getConfig)
|
||||
.selectAs(TeDetectionItemVersion::getConfig, TeQueryTaskOrchestraItemVO::getConfig)
|
||||
.selectAs(TeDetectionItemVersion::getId, TeQueryTaskOrchestraItemVO::getVersionId)
|
||||
.selectAs(TeDetectionItemVersion::getVersionNo, TeQueryTaskOrchestraItemVO::getVersionNo)
|
||||
.leftJoin(TeDetectionItemVersion.class, TeDetectionItemVersion::getId, TeDetectionItem::getPublishedVersionId)
|
||||
.eq(TeTaskOrchestration::getTaskId, taskId)
|
||||
.orderByAsc(TeTaskOrchestration::getOrderNum)
|
||||
.orderByAsc(TeTaskOrchestration::getId);
|
||||
@ -112,10 +116,15 @@ public class TeTaskOrchestrationServiceImpl extends ServiceImpl<TeTaskOrchestrat
|
||||
wrapper
|
||||
.selectAs(TeTaskOrchestration::getItemId, TeQueryTaskDetailDTO::getDetectItemId)
|
||||
.selectAs(TeDetectionItem::getSceneCode, TeQueryTaskDetailDTO::getSceneType)
|
||||
.select(TeDetectionItem::getIsDeploy)
|
||||
.select(TeTaskOrchestration::getOrderNum)
|
||||
.select(TeDetectionItem::getFlowData)
|
||||
.selectAs(TeDetectionItemVersion::getFlowData, TeQueryTaskDetailDTO::getFlowData)
|
||||
.selectAs(TeDetectionItemVersion::getConfig, TeQueryTaskDetailDTO::getConfig)
|
||||
.selectAs(TeDetectionItemVersion::getId, TeQueryTaskDetailDTO::getVersionId)
|
||||
.selectAs(TeDetectionItemVersion::getVersionNo, TeQueryTaskDetailDTO::getVersionNo)
|
||||
.eq(TeTaskOrchestration::getTaskId, taskId)
|
||||
.leftJoin(TeDetectionItem.class, TeDetectionItem::getId, TeTaskOrchestration::getItemId)
|
||||
.leftJoin(TeDetectionItemVersion.class, TeDetectionItemVersion::getId, TeDetectionItem::getPublishedVersionId)
|
||||
.orderByAsc(TeTaskOrchestration::getOrderNum)
|
||||
.orderByAsc(TeTaskOrchestration::getId);
|
||||
|
||||
|
||||
72
sql/workflow_release_version.sql
Normal file
72
sql/workflow_release_version.sql
Normal file
@ -0,0 +1,72 @@
|
||||
-- Drafts remain in te_detection_item. Published versions are immutable snapshots.
|
||||
alter table te_detection_item
|
||||
add column published_version_id varchar(32) null comment 'Current published workflow version' after flow_data,
|
||||
add column published_version_no int null comment 'Current published version number' after published_version_id,
|
||||
add column draft_revision int not null default 0 comment 'Optimistic draft revision' after published_version_no,
|
||||
add column draft_hash varchar(64) null comment 'Draft SHA-256' after draft_revision,
|
||||
add column published_draft_hash varchar(64) null comment 'Last published draft SHA-256' after draft_hash,
|
||||
add index idx_detection_item_published_version (published_version_id);
|
||||
|
||||
create table te_detection_item_version (
|
||||
id varchar(32) not null,
|
||||
detection_item_id varchar(32) not null,
|
||||
version_no int not null,
|
||||
flow_data longtext not null,
|
||||
config longtext null,
|
||||
schema_version int not null default 1,
|
||||
content_hash varchar(64) not null,
|
||||
source_version_id varchar(32) null,
|
||||
publish_note varchar(500) null,
|
||||
create_by varchar(64) null,
|
||||
create_time datetime null,
|
||||
update_by varchar(64) null,
|
||||
update_time datetime null,
|
||||
remark varchar(500) null,
|
||||
primary key (id),
|
||||
unique key uk_detection_item_version (detection_item_id, version_no),
|
||||
key idx_detection_item_version_item (detection_item_id),
|
||||
key idx_detection_item_version_source (source_version_id)
|
||||
) comment='Immutable workflow release history';
|
||||
|
||||
-- Every existing non-empty workflow becomes V1. Older releases did not reliably
|
||||
-- maintain is_deploy, while formal Aima/inspection tasks already executed these definitions.
|
||||
insert into te_detection_item_version
|
||||
(id, detection_item_id, version_no, flow_data, config, schema_version,
|
||||
content_hash, publish_note, create_by, create_time)
|
||||
select replace(uuid(), '-', ''), id, 1, flow_data, config, 1,
|
||||
sha2(concat(ifnull(flow_data, ''), '\n', ifnull(config, '')), 256),
|
||||
'历史流程迁移', create_by, ifnull(update_time, create_time)
|
||||
from te_detection_item
|
||||
where flow_data is not null and flow_data <> '';
|
||||
|
||||
update te_detection_item item
|
||||
join te_detection_item_version version
|
||||
on version.detection_item_id = item.id and version.version_no = 1
|
||||
set item.published_version_id = version.id,
|
||||
item.published_version_no = 1,
|
||||
item.draft_hash = version.content_hash,
|
||||
item.published_draft_hash = version.content_hash,
|
||||
item.is_deploy = '0';
|
||||
|
||||
update te_detection_item
|
||||
set draft_hash = sha2(concat(ifnull(flow_data, ''), '\n', ifnull(config, '')), 256)
|
||||
where draft_hash is null and flow_data is not null;
|
||||
|
||||
create table te_task_inst_flow_version (
|
||||
id varchar(32) not null,
|
||||
inst_id varchar(32) not null,
|
||||
item_id varchar(32) not null,
|
||||
item_order int not null,
|
||||
item_occurrence int not null,
|
||||
version_id varchar(32) not null,
|
||||
version_no int not null,
|
||||
create_by varchar(64) null,
|
||||
create_time datetime null,
|
||||
update_by varchar(64) null,
|
||||
update_time datetime null,
|
||||
remark varchar(500) null,
|
||||
primary key (id),
|
||||
unique key uk_task_inst_flow_item (inst_id, item_order),
|
||||
key idx_task_inst_flow_inst (inst_id),
|
||||
key idx_task_inst_flow_version (version_id)
|
||||
) comment='Workflow release manifest frozen for a task instance';
|
||||
Loading…
Reference in New Issue
Block a user