feat(devops): 实现流水线节点执行功能

- 新增 PipelineJobBean 类作为流水线任务载体
- 重构 PipelineJobManage 类,支持节点级任务管理
- 新增 PipelineNodeJobManage 类用于管理节点任务
- 实现节点任务的启动、中断和删除功能
- 新增节点日志切面和切面类
This commit is contained in:
2025-07-22 23:18:38 +08:00
parent 15038d0c1c
commit f33d14fc77
46 changed files with 1261 additions and 513 deletions
+25
View File
@@ -0,0 +1,25 @@
import request from '@/utils/request'
export function startPipeline() {
return request({
url: 'api/devops/pipeline/start',
method: 'post'
})
}
export function restartPipeline() {
return request({
url: 'api/devops/pipeline/restart',
method: 'post'
})
}
export function queryLastPipelineDetail(instanceId) {
return request({
url: 'api/devops/pipeline/last',
method: 'post',
data: {
instanceId
}
})
}
+15 -1
View File
@@ -75,6 +75,19 @@ const LabelUtil = {
}
}
const SizeUtil = {
formatBytes: function(value) {
if (typeof value !== 'number' || isNaN(value) || value <= 0) {
return '0 Bytes'
}
const units = ['Bytes', 'KB', 'MB', 'GB', 'TB', 'PB', 'EB', 'ZB', 'YB']
const index = Math.floor(Math.log(value) / Math.log(1024))
const unit = units[Math.min(index, units.length - 1)]
const result = (value / Math.pow(1024, index)).toFixed(2)
return `${result} ${unit}`
}
}
const DateUtil = {
formatDate: function(originVal) {
if (originVal === undefined) {
@@ -119,5 +132,6 @@ export {
StringUtil,
ObjectUtil,
LabelUtil,
DateUtil
DateUtil,
SizeUtil
}
@@ -4,13 +4,13 @@
:style="{
borderTopWidth: '5px',
borderTopStyle: 'solid',
borderTopColor: statusColorConst[templateData.status].color
borderTopColor: statusColorConst[pipelineData.currentNodeStatus ? pipelineData.currentNodeStatus : statusColorConst.pending.code].color
}"
>
<div>
<div style="float: left;">
<i
v-if="templateData.currentNodeStatus === 'pending'"
v-if="pipelineData.currentNodeStatus && pipelineData.currentNodeStatus === statusColorConst.pending.code"
class="el-icon-time"
:style="{
fontSize: '22px',
@@ -18,7 +18,7 @@
}"
/>
<i
v-else-if="templateData.currentNodeStatus === 'running'"
v-else-if="pipelineData.currentNodeStatus && pipelineData.currentNodeStatus === statusColorConst.running.code"
class="el-icon-loading"
:style="{
fontSize: '22px',
@@ -26,7 +26,7 @@
}"
/>
<i
v-else-if="templateData.currentNodeStatus === 'success'"
v-else-if="pipelineData.currentNodeStatus && pipelineData.currentNodeStatus === statusColorConst.success.code"
class="el-icon-success"
:style="{
fontSize: '22px',
@@ -34,7 +34,7 @@
}"
/>
<i
v-else-if="templateData.currentNodeStatus === 'fail'"
v-else-if="pipelineData.currentNodeStatus && pipelineData.currentNodeStatus === statusColorConst.fail.code"
class="el-icon-error"
:style="{
fontSize: '22px',
@@ -43,25 +43,25 @@
/>
<i
v-else
class="el-icon-error"
class="el-icon-time"
:style="{
fontSize: '22px',
color: statusColorConst.fail.color
color: statusColorConst.pending.color
}"
/>
</div>
<div class="box-name">{{ templateData.name }}</div>
<div class="box-name">{{ pipelineData.name }}</div>
<div style="clear: both" />
</div>
<div class="box-current-node">
<a @click="onCurrentNodeClick">
{{ templateData.currentNodeMsg }}
{{ pipelineData.currentNodeMsg ? pipelineData.currentNodeMsg : statusColorConst.pending.label }}
</a>
</div>
<el-row>
<el-col :span="12" style="text-align: left">
<!-- 流水线状态为fail,且节点支持重试 -->
<div v-if="templateData.status === 'fail' && templateData.retry === true">
<div v-if="pipelineData.currentNodeStatus === 'fail' && pipelineData.retry === true">
<div class="box-buttons">
<el-button
type="text"
@@ -77,30 +77,13 @@
</el-button>
</div>
</div>
<!-- 流水线状态为running,且有操作按钮 -->
<div v-else-if="templateData.status === 'running' && templateData.buttons && templateData.buttons.length > 0">
<div class="box-buttons">
<el-button
v-for="buttonItem in templateData.buttons"
:key="buttonItem.service"
size="medium"
type="text"
:style="{
padding: 0,
margin: '0 10px 0 10px',
color: ['apply','agree','ok','success'].includes(buttonItem.code) ? 'green' : 'red'
}"
@click="onNodeOperationClick(buttonItem)"
>
{{ buttonItem.title }}
</el-button>
</div>
</div>
<!-- 流水线状态为success,且有需要满足条件的操作按钮 -->
<div v-else-if="templateData.status === 'success' && templateData.buttons && templateData.buttons.length > 0">
<div
v-else-if="pipelineData.currentNodeStatus === statusColorConst.success.code && pipelineData.buttons && pipelineData.buttons.length > 0"
>
<div class="box-buttons">
<el-button
v-for="buttonItem in templateData.buttons"
v-for="buttonItem in pipelineData.buttons"
:key="buttonItem.service"
size="medium"
type="text"
@@ -108,7 +91,28 @@
padding: 0,
margin: '0 10px 0 10px',
color: ['apply','agree','ok','success'].includes(buttonItem.code) ? 'green' : 'red',
display: ['success', 'fail'].includes(buttonItem.code) && templateData.status === buttonItem.code ? '' : 'none'
display: ['success', 'fail'].includes(buttonItem.code) && pipelineData.status === buttonItem.code ? '' : 'none'
}"
@click="onNodeOperationClick(buttonItem)"
>
{{ buttonItem.title }}
</el-button>
</div>
</div>
<!-- 流水线状态非pending,且有操作按钮 -->
<div
v-else-if="pipelineData.currentNodeStatus !== statusColorConst.pending.code && pipelineData.buttons && pipelineData.buttons.length > 0"
>
<div class="box-buttons">
<el-button
v-for="buttonItem in pipelineData.buttons"
:key="buttonItem.service"
size="medium"
type="text"
:style="{
padding: 0,
margin: '0 10px 0 10px',
color: ['apply','agree','ok','success'].includes(buttonItem.code) ? 'green' : 'red'
}"
@click="onNodeOperationClick(buttonItem)"
>
@@ -128,7 +132,7 @@
</el-col>
<div style="clear: both" />
</el-row>
<cute-preview-drawer ref="detailDrawer" :title="templateData.name">
<cute-preview-drawer ref="detailDrawer" :title="pipelineData.name">
<div style="padding: 20px">
这里是明细
</div>
@@ -155,28 +159,26 @@ export default {
data() {
return {
statusColorConst: {
pending: { label: '未开始', color: '#C0C4CC' },
running: { label: '运行中', color: '#67C23A' },
success: { label: '执行成功', color: '#67C23A' },
fail: { label: '执行失败', color: '#F56C6C' }
pending: { code: 'pending', label: '未开始', color: '#C0C4CC' },
running: { code: 'running', label: '运行中', color: '#67C23A' },
success: { code: 'success', label: '执行成功', color: '#67C23A' },
fail: { code: 'fail', label: '执行失败', color: '#F56C6C' }
},
executeTimeStr: ''
executeTimeStr: '',
pipelineData: {}
}
},
watch: {
templateData: {
handler(newVal, oldVal) {
this.pipelineData = { ...newVal }
this.executeTimeStr = this.renderDateTimeStr(newVal)
},
deep: true
}
},
// computed: {
// formatExecuteTimeStr() {
// return this.renderDateTimeStr(this.templateData)
// }
// },
mounted() {
this.pipelineData = this.templateData
this.executeTimeStr = this.renderDateTimeStr(this.templateData)
},
methods: {
@@ -189,7 +191,20 @@ export default {
case 'running':
case 'success':
case 'fail':
executeTimeStr = this.formatTimeDifference(newVal.startTime, newVal.updateTime)
if (newVal.currentNodeStatus !== this.statusColorConst.pending.code) {
if (newVal.startTime && newVal.finishTime) {
executeTimeStr = this.formatTimeDifference(newVal.startTime, newVal.finishTime)
return executeTimeStr
}
if (newVal.startTime) {
executeTimeStr = this.formatTimeDifference(newVal.startTime, new Date())
return executeTimeStr
}
if (newVal.createTime) {
executeTimeStr = this.formatTimeDifference(newVal.createTime, new Date())
return executeTimeStr
}
}
break
default:
executeTimeStr = ''
@@ -223,10 +238,11 @@ export default {
if (seconds >= 0) {
return `${seconds} 秒`
}
return `${diffInMilliseconds} 毫秒`
},
onCurrentNodeClick() {
const templateData = this.templateData || { click: false }
if (!templateData.click) {
const pipelineData = this.pipelineData || { click: false }
if (!pipelineData.click) {
console.warn('当前流水线节点不支持查看明细')
return
}
@@ -234,12 +250,12 @@ export default {
this.$refs.detailDrawer.show()
},
onNodeRetryClick() {
const templateData = this.templateData || { retry: false }
if (!templateData.retry) {
const pipelineData = this.pipelineData || { retry: false }
if (!pipelineData.retry) {
console.warn('当前流水线节点不支持重试')
return
}
console.error('templateData', this.templateData)
console.error('pipelineData', this.pipelineData)
},
onNodeOperationClick(buttonInfo) {
console.error('buttonInfo', buttonInfo)
@@ -277,7 +293,7 @@ export default {
overflow: hidden;
text-overflow: ellipsis;
text-wrap: nowrap;
font-size: 18px;
font-size: 15px;
text-align: center;
}
@@ -1,19 +1,19 @@
<template>
<div>
<el-divider>演示:运行中</el-divider>
<el-divider>演示:动态流水线 <el-button type="primary" @click="startPipelineTest">启动流水线</el-button></el-divider>
<div class="box-pipeline">
<div v-if="dataList1 && dataList1.length > 0" class="box-pipeline-content">
<div v-if="dynamicList && dynamicList.length > 0" class="box-pipeline-content">
<cute-pipeline-node
v-for="(template, index) in dataList1"
v-for="(template, index) in dynamicList"
:key="template.code"
:template-data="dataList1[index]"
:template-data.sync="dynamicList[index]"
/>
</div>
<div v-else class="box-pipeline-content">
<cute-pipeline-node
v-for="(template, index) in templateList"
:key="template.code"
:template-data="dataList1[index]"
:template-data.sync="templateList[index]"
/>
</div>
</div>
@@ -23,14 +23,14 @@
<cute-pipeline-node
v-for="(template, index) in dataList2"
:key="template.code"
:template-data="dataList2[index]"
:template-data.sync="dataList2[index]"
/>
</div>
<div v-else class="box-pipeline-content">
<cute-pipeline-node
v-for="(template, index) in templateList"
:key="template.code"
:template-data="dataList2[index]"
:template-data.sync="templateList[index]"
/>
</div>
</div>
@@ -40,14 +40,14 @@
<cute-pipeline-node
v-for="(template, index) in dataList3"
:key="template.code"
:template-data="dataList3[index]"
:template-data.sync="dataList3[index]"
/>
</div>
<div v-else class="box-pipeline-content">
<cute-pipeline-node
v-for="(template, index) in templateList"
:key="template.code"
:template-data="dataList3[index]"
:template-data.sync="templateList[index]"
/>
</div>
</div>
@@ -58,6 +58,7 @@
import CutePipelineNode from '@/views/components/dev/CutePipelineNode'
import { formatDateTimeStr } from '@/utils'
import { startPipeline, queryLastPipelineDetail } from '@/api/devops/pipeline'
export default {
name: 'CutePipelineNodeDemo',
@@ -154,147 +155,6 @@ export default {
click: false,
retry: true
}],
// 运行中
dataList1: [
{
code: 'node_init',
type: 'service',
name: '初始化',
createTime: '2025-07-17 21:12:01',
updateTime: '2025-07-17 21:13:01',
startTime: '2025-07-17 21:13:01',
currentNodeMsg: '执行成功',
currentNodeStatus: 'success',
status: 'success'
},
{
code: 'node_merge_branch',
type: 'service',
name: '合并代码',
click: true,
retry: true,
createTime: '2025-07-17 21:13:01',
updateTime: '2025-07-17 21:18:01',
startTime: '2025-07-17 21:13:01',
currentNodeMsg: '执行成功',
currentNodeStatus: 'success',
status: 'success'
},
{
code: 'node_build_java',
type: 'service',
name: '构建',
click: true,
retry: true,
detailType: 'gitlab',
parameters: {
pipeline: 'pipeline-backend'
},
createTime: '2025-07-17 21:18:01',
updateTime: '2025-07-17 21:28:01',
startTime: '2025-07-17 21:18:01',
currentNodeMsg: '执行成功',
currentNodeStatus: 'success',
status: 'success'
},
{
code: 'node_image_scan',
type: 'service',
name: '镜像扫描',
click: true,
retry: true,
detailType: 'gitlab',
parameters: {
pipeline: 'pipeline-image-scan'
},
createTime: '2025-07-17 21:28:01',
updateTime: '2025-07-17 21:28:21',
startTime: '2025-07-17 21:28:01',
currentNodeMsg: '执行成功',
currentNodeStatus: 'success',
status: 'success'
},
{
code: 'node_approve_deploy',
type: 'service',
name: '部署审批',
click: false,
retry: true,
buttons: [
{
type: 'execute',
title: '同意',
code: 'agree',
parameters: {}
},
{
type: 'execute',
title: '拒绝',
code: 'refuse',
parameters: {}
}
],
createTime: '2025-07-17 21:28:21',
updateTime: '2025-07-17 21:28:21',
startTime: '2025-07-17 21:28:21',
currentNodeMsg: '等待人工审批',
currentNodeStatus: 'running',
status: 'running'
},
{
code: 'node_deploy_java',
type: 'service',
name: '部署',
click: true,
retry: false,
buttons: [{
type: 'link',
title: '查看部署详情',
code: 'success',
parameters: {
'isBlank': true
}
}],
createTime: '2025-07-17 21:28:21',
updateTime: '2025-07-17 21:28:21',
startTime: null,
currentNodeMsg: '待执行',
currentNodeStatus: 'pending',
status: 'pending'
},
{
code: 'node_merge_confirm',
type: 'service',
name: '合并确认',
click: false,
retry: false,
buttons: [{
type: 'execute',
title: '通过',
code: 'agree',
parameters: {}
}],
createTime: '2025-07-17 21:28:21',
updateTime: '2025-07-17 21:28:21',
startTime: null,
currentNodeMsg: '待执行',
currentNodeStatus: 'pending',
status: 'pending'
},
{
code: 'node_merge_master',
type: 'service',
name: '合并回Master',
click: false,
retry: true,
createTime: '2025-07-17 21:28:21',
updateTime: '2025-07-17 21:28:21',
startTime: null,
currentNodeMsg: '待执行',
currentNodeStatus: 'pending',
status: 'pending'
}
],
// 成功
dataList2: [
{
@@ -576,21 +436,18 @@ export default {
currentNodeStatus: 'pending',
status: 'pending'
}
]
],
// 动态流水线
dynamicHook: null,
dynamicList: [],
dynamicInstanceId: ''
}
},
mounted() {
setInterval(this.refreshData, 1000)
this.refreshData()
},
methods: {
refreshData() {
// 运行中
for (const datum of this.dataList1) {
if (datum.status === 'running') {
datum.updateTime = formatDateTimeStr(new Date())
}
}
this.dataList1 = [...this.dataList1]
// 成功
for (const datum of this.dataList2) {
if (datum.status === 'running') {
@@ -605,6 +462,22 @@ export default {
}
}
this.dataList3 = [...this.dataList3]
},
startPipelineTest() {
const that = this
if (that.dynamicHook) {
clearInterval(that.dynamicHook)
}
startPipeline().then((instanceId) => {
that.dynamicInstanceId = instanceId
setTimeout(() => {
that.dynamicHook = setInterval(() => {
queryLastPipelineDetail(that.dynamicInstanceId).then((data) => {
that.dynamicList = data
})
}, 1000)
}, 2000)
})
}
}
}
@@ -22,10 +22,13 @@ public interface PipelineConst {
// String INSTANCE_NODE_TEMPLATE_ARGS_KEY = "pipelineInstanceNodeTemplateArgs";
// String INSTANCE_RETRY_NODE_INDEX_KEY = "pipelineInstanceRetryNodeIndex";
// 以下为标准常量
String INSTANCE = "Instance";
String INSTANCE_ID = "InstanceId";
String TEMPLATE = "Template";
String TEMPLATE_ID = "TemplateId";
String ENV = "Env";
String CONTEXT_NAME = "ContextName";
String RETRY_NODE_CODE = "RetryNodeCode";
String CURRENT_NODE_TEMPLATE = "CurrentNodeTemplate";
String LAST_NODE_RESULT = "LastNodeResult";
}
@@ -3,27 +3,37 @@ package cn.odboy.devops.controller;
import cn.odboy.annotation.AnonymousGetMapping;
import cn.odboy.devops.dal.dataobject.PipelineInstanceTb;
import cn.odboy.devops.dal.dataobject.PipelineTemplateTb;
import cn.odboy.devops.dal.model.DevOpsQueryLastPipelineDetailArgs;
import cn.odboy.devops.service.PipelineInstanceService;
import cn.odboy.devops.service.PipelineTemplateService;
import io.swagger.annotations.Api;
import io.swagger.annotations.ApiOperation;
import lombok.RequiredArgsConstructor;
import lombok.extern.slf4j.Slf4j;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.http.ResponseEntity;
import org.springframework.security.access.prepost.PreAuthorize;
import org.springframework.validation.annotation.Validated;
import org.springframework.web.bind.annotation.PostMapping;
import org.springframework.web.bind.annotation.RequestBody;
import org.springframework.web.bind.annotation.RequestMapping;
import org.springframework.web.bind.annotation.RestController;
@Slf4j
@RestController
@RequiredArgsConstructor
@RequestMapping("/api/devops/pipeline")
@Api(tags = "系统:验证码管理")
@Api(tags = "DevOps:流水线管理")
public class DevopsPipelineController {
@Autowired
private PipelineTemplateService pipelineTemplateService;
@Autowired
private PipelineInstanceService pipelineInstanceService;
@AnonymousGetMapping(value = "/start")
public ResponseEntity<?> start() {
@ApiOperation("启动流水线")
@PostMapping(value = "/start")
@PreAuthorize("@el.check()")
public ResponseEntity<?> startPipeline() {
// 输入参数
String appName = "cuteops";
String envCode = "daily";
@@ -37,7 +47,33 @@ public class DevopsPipelineController {
pipelineInstanceTb.setContextName(appName);
pipelineInstanceTb.setEnv(envCode);
pipelineInstanceTb.setTemplateContent(pipelineTemplateTb.getTemplate());
pipelineInstanceService.start(pipelineInstanceTb);
return ResponseEntity.ok("流水线启动成功");
return ResponseEntity.ok(pipelineInstanceService.startPipeline(pipelineInstanceTb));
}
@ApiOperation("重启流水线")
@PostMapping(value = "/restart")
@PreAuthorize("@el.check()")
public ResponseEntity<?> restartPipeline() {
String appName = "cuteops";
String envCode = "daily";
PipelineTemplateTb pipelineTemplateTb = pipelineTemplateService.getPipelineTemplateById(4L);
PipelineInstanceTb pipelineInstanceTb = new PipelineInstanceTb();
pipelineInstanceTb.setTemplateId(pipelineTemplateTb.getId());
pipelineInstanceTb.setInstanceName("流水线测试");
pipelineInstanceTb.setTemplateType(pipelineTemplateTb.getType());
pipelineInstanceTb.setContextName(appName);
pipelineInstanceTb.setEnv(envCode);
pipelineInstanceTb.setTemplateContent(pipelineTemplateTb.getTemplate());
pipelineInstanceTb.setInstanceId(1947604330727079936L);
String retryNodeCode = "node_build_java";
pipelineInstanceService.restartPipeline(pipelineInstanceTb, retryNodeCode);
return ResponseEntity.ok("流水线重启成功");
}
@ApiOperation("查询流水线明细")
@PostMapping(value = "/last")
@PreAuthorize("@el.check()")
public ResponseEntity<?> queryLastPipelineDetail(@Validated @RequestBody DevOpsQueryLastPipelineDetailArgs args) {
return ResponseEntity.ok(pipelineInstanceService.queryLastPipelineDetail(args.getInstanceId()));
}
}
@@ -32,7 +32,7 @@ public class PipelineInstanceNodeDetailTb extends CsObject {
* 流水线节点Id
*/
@CollectionField("node_id")
private String nodeId;
private Long nodeId;
/**
* 开始时间
@@ -46,11 +46,6 @@ public class PipelineInstanceNodeDetailTb extends CsObject {
@CollectionField("finish_time")
private Date finishTime;
/**
* 节点索引
*/
@CollectionField("node_index")
private Integer nodeIndex;
/**
* 节点编码
*/
@@ -72,6 +72,8 @@ public class PipelineInstanceNodeTb extends CsObject {
private Date updateTime;
@CollectionField("start_time")
private Date startTime;
@CollectionField("finish_time")
private Date finishTime;
/**
* 节点状态
*/
@@ -79,9 +81,4 @@ public class PipelineInstanceNodeTb extends CsObject {
private String currentNodeMsg = PipelineStatusEnum.PENDING.getDesc();
@CollectionField("current_node_status")
private String currentNodeStatus = PipelineStatusEnum.PENDING.getCode();
/**
* 节点状态
*/
@CollectionField("status")
private String status = PipelineStatusEnum.PENDING.getCode();
}
@@ -0,0 +1,14 @@
package cn.odboy.devops.dal.model;
import cn.odboy.base.CsObject;
import lombok.Getter;
import lombok.Setter;
import javax.validation.constraints.NotBlank;
@Getter
@Setter
public class DevOpsQueryLastPipelineDetailArgs extends CsObject {
@NotBlank(message = "流水线实例Id必填")
private String instanceId;
}
@@ -20,8 +20,8 @@ public class PipelineInstanceDAO {
if (o != null) {
return true;
}
// 365天
redisHelper.set(redisKey, true, 60 * 60 * 24 * 365);
// 3个月
redisHelper.set(redisKey, true, 60 * 60 * 24 * 90);
return false;
}
@@ -0,0 +1,30 @@
package cn.odboy.devops.framework.pipeline;
import cn.odboy.devops.constant.pipeline.PipelineStatusEnum;
import cn.odboy.devops.dal.dataobject.PipelineInstanceNodeTb;
import cn.odboy.devops.service.PipelineInstanceNodeDetailService;
import cn.odboy.devops.service.PipelineInstanceNodeService;
import org.springframework.beans.factory.annotation.Autowired;
import java.util.Date;
public abstract class AbstractPipelineNodeJobService {
// @Autowired
// private PipelineTemplateService pipelineTemplateService;
// @Autowired
// private PipelineInstanceService pipelineInstanceService;
@Autowired
private PipelineInstanceNodeService pipelineInstanceNodeService;
@Autowired
private PipelineInstanceNodeDetailService pipelineInstanceNodeDetailService;
public PipelineInstanceNodeTb getCurrentNodeInfo(long instanceId, String code) {
return pipelineInstanceNodeService.getPipelineInstanceNodeByArgs(instanceId, code);
}
public void addLog(PipelineInstanceNodeTb pipelineInstanceNode, String stepName, PipelineStatusEnum stepStatus, String stepMsg, Date finishTime) {
if (pipelineInstanceNode != null) {
pipelineInstanceNodeDetailService.addLog(pipelineInstanceNode.getId(), pipelineInstanceNode.getCode(), stepName, stepStatus, stepMsg, finishTime);
}
}
}
@@ -0,0 +1,205 @@
package cn.odboy.devops.framework.pipeline.core;
import cn.hutool.core.thread.ThreadUtil;
import cn.hutool.core.util.StrUtil;
import cn.odboy.devops.constant.pipeline.PipelineConst;
import cn.odboy.devops.constant.pipeline.PipelineStatusEnum;
import cn.odboy.devops.dal.dataobject.PipelineInstanceNodeTb;
import cn.odboy.devops.dal.dataobject.PipelineInstanceTb;
import cn.odboy.devops.dal.mysql.PipelineInstanceMapper;
import cn.odboy.devops.dal.redis.PipelineInstanceDAO;
import cn.odboy.devops.framework.pipeline.model.PipelineNodeTemplateVo;
import cn.odboy.devops.service.PipelineInstanceNodeService;
import cn.odboy.framework.context.SpringBeanHolder;
import cn.odboy.framework.exception.BadRequestException;
import com.alibaba.fastjson2.JSON;
import lombok.extern.slf4j.Slf4j;
import org.quartz.*;
import java.util.List;
/**
* 流水线任务载体
*
* @author odboy
* @date 2025-07-19
*/
@Slf4j
public class PipelineJobBean implements InterruptableJob {
private volatile boolean interrupted = false;
private volatile Long currentInstanceId = null;
private volatile String currentNodeCode = null;
@Override
public void execute(JobExecutionContext jobExecutionContext) throws JobExecutionException {
// 获取服务
PipelineInstanceMapper pipelineInstanceMapper = SpringBeanHolder.getBean(PipelineInstanceMapper.class);
PipelineInstanceNodeService pipelineInstanceNodeService = SpringBeanHolder.getBean(PipelineInstanceNodeService.class);
PipelineInstanceDAO pipelineInstanceDAO = SpringBeanHolder.getBean(PipelineInstanceDAO.class);
PipelineJobManage pipelineJobManage = SpringBeanHolder.getBean(PipelineJobManage.class);
PipelineNodeJobManage pipelineNodeJobManage = SpringBeanHolder.getBean(PipelineNodeJobManage.class);
// 参数解析
JobDataMap jobDataMap = jobExecutionContext.getMergedJobDataMap();
String retryNodeCode = jobDataMap.getString(PipelineConst.RETRY_NODE_CODE);
PipelineInstanceTb pipelineInstanceTb = (PipelineInstanceTb) jobDataMap.get(PipelineConst.INSTANCE);
pipelineInstanceTb.setStatus(PipelineStatusEnum.PENDING.getCode());
// // 模拟终止流水线
// ThreadUtil.execAsync(() -> {
// ThreadUtil.safeSleep(6000);
// try {
// pipelineJobManage.interruptJob(pipelineInstanceTb.getInstanceId());
// } catch (SchedulerException e) {
// log.error("终止流水线异常", e);
// }
// });
if (StrUtil.isBlank(retryNodeCode)) {
// 创建实例
pipelineInstanceMapper.insert(pipelineInstanceTb);
// 创建实例明细
pipelineInstanceNodeService.createPipelineInstanceDetail(pipelineInstanceTb);
} else {
// 重建实例明细
pipelineInstanceNodeService.remakePipelineInstanceDetail(pipelineInstanceTb, retryNodeCode);
}
// 实例节点模板
String templateContent = pipelineInstanceTb.getTemplateContent();
List<PipelineNodeTemplateVo> pipelineNodeTemplateVos = JSON.parseArray(templateContent, PipelineNodeTemplateVo.class);
Long instanceId = pipelineInstanceTb.getInstanceId();
// 实例状态(默认是成功的)
PipelineStatusEnum pipelineStatusEnum = PipelineStatusEnum.SUCCESS;
try {
// 更新实例状态(running)
pipelineInstanceTb.setStatus(PipelineStatusEnum.RUNNING.getCode());
pipelineInstanceMapper.updateById(pipelineInstanceTb);
// 依次执行流水线节点任务
boolean skipFlag = true;
boolean foundRetryNode = false;
this.currentInstanceId = pipelineInstanceTb.getInstanceId();
for (PipelineNodeTemplateVo pipelineNodeTemplateVo : pipelineNodeTemplateVos) {
this.currentNodeCode = pipelineNodeTemplateVo.getCode();
// 如果 retryNodeCode 不为空,进入跳过逻辑
if (StrUtil.isNotBlank(retryNodeCode)) {
if (skipFlag && !pipelineNodeTemplateVo.getCode().equals(retryNodeCode)) {
// 如果当前节点不是重试节点,且 skipFlag 为 true,则跳过
continue;
} else {
// 找到重试节点,开始执行
skipFlag = false;
foundRetryNode = true;
}
}
if (StrUtil.isNotBlank(retryNodeCode) && !foundRetryNode) {
log.error("未找到重试节点: {}", retryNodeCode);
throw new BadRequestException("未找到重试节点: " + retryNodeCode);
}
try {
pipelineNodeJobManage.startJob(pipelineInstanceTb, pipelineNodeTemplateVo);
pipelineInstanceTb.setCurrentNode(pipelineNodeTemplateVo.getCode());
pipelineInstanceTb.setCurrentNodeStatus(PipelineStatusEnum.RUNNING.getCode());
pipelineInstanceMapper.updateById(pipelineInstanceTb);
// 监听节点执行状态
while (true) {
// 检查中断状态
// if (Thread.currentThread().isInterrupted() || interrupted) {
if (interrupted) {
log.error("检测到中断,停止流水线执行");
pipelineStatusEnum = PipelineStatusEnum.FAIL;
pipelineInstanceTb.setCurrentNode(pipelineNodeTemplateVo.getCode());
pipelineInstanceTb.setCurrentNodeStatus(pipelineStatusEnum.getCode());
pipelineInstanceMapper.updateById(pipelineInstanceTb);
break;
}
ThreadUtil.safeSleep(5000);
PipelineInstanceNodeTb pipelineInstanceNode = pipelineInstanceNodeService.getPipelineInstanceNodeByArgs(instanceId, pipelineNodeTemplateVo.getCode());
if (PipelineStatusEnum.SUCCESS.getCode().equals(pipelineInstanceNode.getCurrentNodeStatus())) {
pipelineStatusEnum = PipelineStatusEnum.SUCCESS;
pipelineInstanceTb.setCurrentNodeStatus(pipelineStatusEnum.getCode());
pipelineInstanceMapper.updateById(pipelineInstanceTb);
break;
}
if (PipelineStatusEnum.FAIL.getCode().equals(pipelineInstanceNode.getCurrentNodeStatus())) {
pipelineStatusEnum = PipelineStatusEnum.FAIL;
pipelineInstanceTb.setCurrentNodeStatus(pipelineStatusEnum.getCode());
pipelineInstanceMapper.updateById(pipelineInstanceTb);
break;
}
}
if (pipelineStatusEnum.equals(PipelineStatusEnum.FAIL)) {
// 节点失败,全失败
break;
}
} catch (Exception e) {
log.error("流水线节点运行失败", e);
pipelineInstanceTb.setStatus(PipelineStatusEnum.FAIL.getCode());
pipelineInstanceMapper.updateById(pipelineInstanceTb);
}
}
// 更新实例状态
pipelineInstanceTb.setStatus(pipelineStatusEnum.getCode());
pipelineInstanceMapper.updateById(pipelineInstanceTb);
pipelineInstanceDAO.unLock(pipelineInstanceTb);
} catch (Exception e) {
log.error("流水线运行失败", e);
pipelineInstanceTb.setStatus(PipelineStatusEnum.FAIL.getCode());
pipelineInstanceMapper.updateById(pipelineInstanceTb);
pipelineInstanceDAO.unLock(pipelineInstanceTb);
} finally {
pipelineJobManage.deleteJob(instanceId);
}
}
/**
* 流水线任务已中断
*/
@Override
public void interrupt() throws UnableToInterruptJobException {
this.interrupted = true;
if (this.currentInstanceId != null && StrUtil.isNotBlank(this.currentNodeCode)) {
PipelineInstanceMapper pipelineInstanceMapper = SpringBeanHolder.getBean(PipelineInstanceMapper.class);
PipelineInstanceDAO pipelineInstanceDAO = SpringBeanHolder.getBean(PipelineInstanceDAO.class);
int seconds = 30_000;
log.info("中断信号已发出,等待最多 {} 秒,若未自行结束则强制终止线程", seconds);
long waitStart = System.currentTimeMillis();
while ((System.currentTimeMillis() - waitStart) < seconds) {
if (Thread.currentThread().isInterrupted()) {
log.info("线程已自行中断");
return;
}
ThreadUtil.safeSleep(1000);
}
log.info("等待超时,强制终止线程");
PipelineInstanceTb pipelineInstanceTb = pipelineInstanceMapper.selectById(this.currentInstanceId);
if (pipelineInstanceTb != null) {
// 更新状态 & 解锁
pipelineInstanceTb.setCurrentNode(this.currentNodeCode);
pipelineInstanceTb.setCurrentNodeStatus(PipelineStatusEnum.FAIL.getCode());
pipelineInstanceTb.setStatus(PipelineStatusEnum.FAIL.getCode());
pipelineInstanceMapper.updateById(pipelineInstanceTb);
pipelineInstanceDAO.unLock(pipelineInstanceTb);
// 杀死
try {
log.info("流水线实例线程被杀死, instanceId={}", this.currentInstanceId);
Thread.currentThread().stop();
} catch (Exception e) {
log.error("流水线实例线程杀死失败", e);
}
}
}
}
}
@@ -1,15 +1,14 @@
package cn.odboy.devops.framework.pipeline.core;
import cn.hutool.core.util.IdUtil;
import cn.odboy.devops.constant.pipeline.PipelineConst;
import cn.odboy.devops.dal.dataobject.PipelineInstanceTb;
import cn.odboy.devops.framework.pipeline.model.PipelineNodeTemplateVo;
import cn.odboy.framework.exception.BadRequestException;
import lombok.extern.slf4j.Slf4j;
import org.quartz.*;
import org.springframework.stereotype.Component;
import javax.annotation.Resource;
import javax.validation.constraints.NotNull;
/**
* 流水线任务管理器
@@ -20,37 +19,29 @@ import javax.validation.constraints.NotNull;
@Slf4j
@Component
public class PipelineJobManage {
private static final String JOB_KEY = "PIPELINE_JOB_%s_%s";
private static final String JOB_KEY = "PIPELINE_JOB_%s";
@Resource
private Scheduler scheduler;
/**
* 启动 job
*/
public void startJob(PipelineInstanceTb pipelineInstance, PipelineNodeTemplateVo pipelineNodeTemplate) {
public String startJob(PipelineInstanceTb pipelineInstance) {
// 流水线Id(雪花)
Long instanceId = IdUtil.getSnowflakeNextId();
try {
// 流水线Id
Long instanceId = pipelineInstance.getInstanceId();
String jobId = String.format(JOB_KEY, instanceId, pipelineNodeTemplate.getCode());
pipelineInstance.setInstanceId(instanceId);
String keyName = String.format(JOB_KEY, instanceId);
// 构建 JobDetail
JobKey jobKey = JobKey.jobKey(jobId);
JobDetail jobDetail = JobBuilder
.newJob(PipelineNodeJobBean.class)
.withIdentity(jobKey)
.build();
JobKey jobKey = JobKey.jobKey(keyName);
JobDetail jobDetail = JobBuilder.newJob(PipelineJobBean.class).withIdentity(jobKey).build();
// 构建Trigger
TriggerKey triggerKey = TriggerKey.triggerKey(String.format(JOB_KEY, instanceId, pipelineNodeTemplate.getCode()));
Trigger cronTrigger = TriggerBuilder.newTrigger()
.withIdentity(triggerKey)
.startNow()
.build();
TriggerKey triggerKey = TriggerKey.triggerKey(keyName);
Trigger cronTrigger = TriggerBuilder.newTrigger().withIdentity(triggerKey).startNow().build();
// 任务参数
JobDataMap jobDataMap = cronTrigger.getJobDataMap();
jobDataMap.put(PipelineConst.INSTANCE_ID, pipelineInstance.getInstanceId());
jobDataMap.put(PipelineConst.CONTEXT_NAME, pipelineInstance.getContextName());
jobDataMap.put(PipelineConst.ENV, pipelineInstance.getEnv());
jobDataMap.put(PipelineConst.TEMPLATE_ID, pipelineInstance.getTemplateId());
jobDataMap.put(PipelineConst.CURRENT_NODE_TEMPLATE, pipelineNodeTemplate);
jobDataMap.put(PipelineConst.INSTANCE, pipelineInstance);
jobDataMap.put(PipelineConst.RETRY_NODE_CODE, "");
try {
// 在quartz的内部线程池中异步执行
scheduler.scheduleJob(jobDetail, cronTrigger);
@@ -61,14 +52,15 @@ public class PipelineJobManage {
log.error("创建定时任务失败", e);
throw new BadRequestException("创建定时任务失败");
}
return String.valueOf(instanceId);
}
/**
* 删除job
*/
public void deleteJob(@NotNull String instanceId, @NotNull String nodeCode) {
public void deleteJob(Long instanceId) {
try {
JobKey jobKey = JobKey.jobKey(String.format(JOB_KEY, instanceId, nodeCode));
JobKey jobKey = JobKey.jobKey(String.format(JOB_KEY, instanceId));
scheduler.pauseJob(jobKey);
scheduler.deleteJob(jobKey);
} catch (Exception e) {
@@ -81,8 +73,34 @@ public class PipelineJobManage {
* 中断正在执行的任务<br/>
* 任务终止需要配合响应中断或停止信号
*/
public void interruptJob(@NotNull String instanceId, @NotNull String nodeCode) throws SchedulerException {
JobKey jobKey = JobKey.jobKey(String.format(JOB_KEY, instanceId, nodeCode));
public void interruptJob(Long instanceId) throws SchedulerException {
JobKey jobKey = JobKey.jobKey(String.format(JOB_KEY, instanceId));
scheduler.interrupt(jobKey);
}
public void startJobByNodeCode(PipelineInstanceTb pipelineInstance, String retryNodeCode) {
try {
Long instanceId = pipelineInstance.getInstanceId();
String keyName = String.format(JOB_KEY, instanceId);
// 构建 JobDetail
JobKey jobKey = JobKey.jobKey(keyName);
JobDetail jobDetail = JobBuilder.newJob(PipelineJobBean.class).withIdentity(jobKey).build();
// 构建Trigger
TriggerKey triggerKey = TriggerKey.triggerKey(keyName);
Trigger cronTrigger = TriggerBuilder.newTrigger().withIdentity(triggerKey).startNow().build();
// 任务参数
JobDataMap jobDataMap = cronTrigger.getJobDataMap();
jobDataMap.put(PipelineConst.INSTANCE, pipelineInstance);
jobDataMap.put(PipelineConst.RETRY_NODE_CODE, retryNodeCode);
try {
// 在quartz的内部线程池中异步执行
scheduler.scheduleJob(jobDetail, cronTrigger);
} catch (ObjectAlreadyExistsException e) {
log.warn("定时任务已存在,跳过加载");
}
} catch (Exception e) {
log.error("创建定时任务失败", e);
throw new BadRequestException("创建定时任务失败");
}
}
}
@@ -1,10 +1,19 @@
package cn.odboy.devops.framework.pipeline.core;
import cn.odboy.devops.constant.pipeline.PipelineConst;
import cn.odboy.devops.constant.pipeline.PipelineStatusEnum;
import cn.odboy.devops.framework.pipeline.model.PipelineNodeJobExecuteResult;
import cn.odboy.devops.framework.pipeline.model.PipelineNodeTemplateVo;
import cn.odboy.devops.service.PipelineInstanceNodeService;
import cn.odboy.framework.context.SpringBeanHolder;
import cn.odboy.framework.exception.BadRequestException;
import lombok.extern.slf4j.Slf4j;
import org.quartz.*;
import org.quartz.Job;
import org.quartz.JobDataMap;
import org.quartz.JobExecutionContext;
import org.quartz.JobExecutionException;
import java.util.Date;
/**
* 流水线任务载体
@@ -13,18 +22,29 @@ import org.quartz.*;
* @date 2025-07-19
*/
@Slf4j
public class PipelineNodeJobBean implements InterruptableJob {
public class PipelineNodeJobBean implements Job {
@Override
public void execute(JobExecutionContext jobExecutionContext) throws JobExecutionException {
PipelineInstanceNodeService pipelineInstanceNodeService = SpringBeanHolder.getBean(PipelineInstanceNodeService.class);
JobDataMap jobDataMap = jobExecutionContext.getMergedJobDataMap();
long instanceId = jobDataMap.getLong(PipelineConst.INSTANCE_ID);
PipelineNodeTemplateVo currentNodeTemplate = (PipelineNodeTemplateVo) jobDataMap.get(PipelineConst.CURRENT_NODE_TEMPLATE);
String serviceName = PipelineConst.EXECUTOR_PREFIX + currentNodeTemplate.getCode();
String nodeCode = currentNodeTemplate.getCode();
try {
String serviceName = PipelineConst.EXECUTOR_PREFIX + nodeCode;
PipelineNodeJobExecutor pipelineNodeJobExecutor = SpringBeanHolder.getBean(serviceName);
pipelineNodeJobExecutor.execute(jobDataMap);
pipelineInstanceNodeService.updatePipelineInstanceNodeByArgs(instanceId, nodeCode, PipelineStatusEnum.RUNNING, PipelineStatusEnum.RUNNING.getDesc());
PipelineNodeJobExecuteResult executeResult = pipelineNodeJobExecutor.execute(jobDataMap);
jobExecutionContext.getMergedJobDataMap().put(PipelineConst.LAST_NODE_RESULT, executeResult);
pipelineInstanceNodeService.finishPipelineInstanceNodeByArgs(instanceId, nodeCode, PipelineStatusEnum.SUCCESS, PipelineStatusEnum.SUCCESS.getDesc(), new Date());
} catch (BadRequestException e) {
log.error("流水线节点执行异常", e);
pipelineInstanceNodeService.finishPipelineInstanceNodeByArgs(instanceId, nodeCode, PipelineStatusEnum.FAIL, e.getMessage(), new Date());
} catch (Exception e) {
log.error("流水线节点执行异常", e);
pipelineInstanceNodeService.finishPipelineInstanceNodeByArgs(instanceId, nodeCode, PipelineStatusEnum.FAIL, PipelineStatusEnum.FAIL.getDesc(), new Date());
}
@Override
public void interrupt() throws UnableToInterruptJobException {
}
}
@@ -1,6 +1,7 @@
package cn.odboy.devops.framework.pipeline.core;
import cn.odboy.devops.framework.pipeline.model.PipelineNodeJobExecuteResult;
import cn.odboy.framework.exception.BadRequestException;
import org.quartz.JobDataMap;
/**
@@ -10,5 +11,5 @@ import org.quartz.JobDataMap;
* @date 2025-07-21
*/
public interface PipelineNodeJobExecutor {
PipelineNodeJobExecuteResult execute(JobDataMap contextArgs);
PipelineNodeJobExecuteResult execute(JobDataMap contextArgs) throws BadRequestException;
}
@@ -0,0 +1,89 @@
package cn.odboy.devops.framework.pipeline.core;
import cn.odboy.devops.constant.pipeline.PipelineConst;
import cn.odboy.devops.dal.dataobject.PipelineInstanceTb;
import cn.odboy.devops.framework.pipeline.model.PipelineNodeTemplateVo;
import cn.odboy.framework.exception.BadRequestException;
import lombok.extern.slf4j.Slf4j;
import org.quartz.*;
import org.springframework.stereotype.Component;
import javax.annotation.Resource;
import javax.validation.constraints.NotNull;
/**
* 流水线任务管理器
*
* @author odboy
* @date 2025-07-19
*/
@Slf4j
@Component
public class PipelineNodeJobManage {
private static final String JOB_KEY = "PIPELINE_NODE_JOB_%s_%s";
@Resource
private Scheduler scheduler;
/**
* 启动 job
*/
public void startJob(PipelineInstanceTb pipelineInstance, PipelineNodeTemplateVo pipelineNodeTemplate) {
try {
// 流水线Id
Long instanceId = pipelineInstance.getInstanceId();
String keyName = String.format(JOB_KEY, instanceId, pipelineNodeTemplate.getCode());
// 构建 JobDetail
JobKey jobKey = JobKey.jobKey(keyName);
JobDetail jobDetail = JobBuilder
.newJob(PipelineNodeJobBean.class)
.withIdentity(jobKey)
.build();
// 构建Trigger
TriggerKey triggerKey = TriggerKey.triggerKey(keyName);
Trigger cronTrigger = TriggerBuilder.newTrigger()
.withIdentity(triggerKey)
.startNow()
.build();
// 任务参数
JobDataMap jobDataMap = cronTrigger.getJobDataMap();
jobDataMap.put(PipelineConst.INSTANCE_ID, pipelineInstance.getInstanceId());
jobDataMap.put(PipelineConst.CONTEXT_NAME, pipelineInstance.getContextName());
jobDataMap.put(PipelineConst.ENV, pipelineInstance.getEnv());
jobDataMap.put(PipelineConst.TEMPLATE, pipelineInstance.getTemplateContent());
jobDataMap.put(PipelineConst.TEMPLATE_ID, pipelineInstance.getTemplateId());
jobDataMap.put(PipelineConst.CURRENT_NODE_TEMPLATE, pipelineNodeTemplate);
try {
// 在quartz的内部线程池中异步执行
scheduler.scheduleJob(jobDetail, cronTrigger);
} catch (ObjectAlreadyExistsException e) {
log.warn("定时任务已存在,跳过加载");
}
} catch (Exception e) {
log.error("创建定时任务失败", e);
throw new BadRequestException("创建定时任务失败");
}
}
/**
* 删除job
*/
public void deleteJob(@NotNull String instanceId, @NotNull String nodeCode) {
try {
JobKey jobKey = JobKey.jobKey(String.format(JOB_KEY, instanceId, nodeCode));
scheduler.pauseJob(jobKey);
scheduler.deleteJob(jobKey);
} catch (Exception e) {
log.error("删除定时任务失败", e);
throw new BadRequestException("删除定时任务失败");
}
}
/**
* 中断正在执行的任务<br/>
* 任务终止需要配合响应中断或停止信号
*/
public void interruptJob(@NotNull String instanceId, @NotNull String nodeCode) throws SchedulerException {
JobKey jobKey = JobKey.jobKey(String.format(JOB_KEY, instanceId, nodeCode));
scheduler.interrupt(jobKey);
}
}
@@ -0,0 +1,18 @@
package cn.odboy.devops.framework.pipeline.log;
import java.lang.annotation.ElementType;
import java.lang.annotation.Retention;
import java.lang.annotation.RetentionPolicy;
import java.lang.annotation.Target;
/**
* 流水线节点明细日志 切面
*
* @author odboy
* @date 2025-07-22
*/
@Retention(RetentionPolicy.RUNTIME)
@Target(ElementType.METHOD)
public @interface PipelineNodeStepLog {
String value() default "默认步骤名";
}
@@ -0,0 +1,63 @@
package cn.odboy.devops.framework.pipeline.log;
import cn.odboy.devops.constant.pipeline.PipelineStatusEnum;
import cn.odboy.devops.dal.dataobject.PipelineInstanceNodeTb;
import cn.odboy.devops.framework.pipeline.AbstractPipelineNodeJobService;
import cn.odboy.framework.exception.BadRequestException;
import lombok.extern.slf4j.Slf4j;
import org.aspectj.lang.ProceedingJoinPoint;
import org.aspectj.lang.annotation.Around;
import org.aspectj.lang.annotation.Aspect;
import org.aspectj.lang.annotation.Pointcut;
import org.aspectj.lang.reflect.MethodSignature;
import org.springframework.stereotype.Component;
import java.lang.reflect.Method;
import java.util.Date;
@Slf4j
@Aspect
@Component
public class PipelineNodeStepLogAspect extends AbstractPipelineNodeJobService {
/**
* 配置切入点
*/
@Pointcut("@annotation(cn.odboy.devops.framework.pipeline.log.PipelineNodeStepLog)")
public void logPointcut() {
// 该方法无方法体,主要为了让同类中其他方法使用此切入点
}
/**
* 配置环绕通知,使用在方法logPointcut()上注册的切入点
*
* @param joinPoint join point for advice
*/
@Around("logPointcut()")
public Object logAround(ProceedingJoinPoint joinPoint) throws Throwable {
MethodSignature signature = (MethodSignature) joinPoint.getSignature();
Method method = signature.getMethod();
PipelineInstanceNodeTb currentNodeInfo = null;
for (Object arg : joinPoint.getArgs()) {
if (arg instanceof PipelineInstanceNodeTb) {
currentNodeInfo = (PipelineInstanceNodeTb) arg;
break;
}
}
PipelineNodeStepLog pipelineNodeStepLog = method.getAnnotation(PipelineNodeStepLog.class);
try {
addLog(currentNodeInfo, pipelineNodeStepLog.value(), PipelineStatusEnum.RUNNING, PipelineStatusEnum.RUNNING.getDesc(), null);
Object result = joinPoint.proceed();
addLog(currentNodeInfo, pipelineNodeStepLog.value(), PipelineStatusEnum.SUCCESS, PipelineStatusEnum.SUCCESS.getDesc(), new Date());
return result;
} catch (BadRequestException e) {
log.error("执行失败", e);
addLog(currentNodeInfo, pipelineNodeStepLog.value(), PipelineStatusEnum.FAIL, e.getMessage(), new Date());
throw e;
} catch (Throwable e) {
log.error("执行失败", e);
addLog(currentNodeInfo, pipelineNodeStepLog.value(), PipelineStatusEnum.FAIL, PipelineStatusEnum.FAIL.getDesc(), new Date());
throw e;
}
}
}
@@ -0,0 +1,44 @@
package cn.odboy.devops.framework.pipeline.model;
import lombok.Getter;
import lombok.Setter;
import java.util.Date;
/**
* 流水线节点数据
*
* @author odboy
*/
@Getter
@Setter
public class PipelineNodeDataVo extends PipelineNodeTemplateVo {
/**
* 节点创建时间
*/
private Date createTime;
/**
* 节点更新时间
*/
private Date updateTime;
/**
* 节点开始运行时间
*/
private Date startTime;
/**
* 节点结束运行时间
*/
private Date finishTime;
/**
* 流水线节点步骤信息
*/
private String currentNodeMsg;
/**
* 流水线节点状态
*/
private String currentNodeStatus;
/**
* 流水线实例状态
*/
private String status;
}
@@ -10,18 +10,29 @@ import lombok.Setter;
public class PipelineNodeJobExecuteResult extends CsObject {
private PipelineStatusEnum status;
private String message;
private Object data;
public static PipelineNodeJobExecuteResult success() {
PipelineNodeJobExecuteResult result = new PipelineNodeJobExecuteResult();
result.setStatus(PipelineStatusEnum.SUCCESS);
result.setMessage(PipelineStatusEnum.SUCCESS.getDesc());
result.setData("");
return result;
}
public static PipelineNodeJobExecuteResult success(String message) {
public static PipelineNodeJobExecuteResult success(Object data) {
PipelineNodeJobExecuteResult result = new PipelineNodeJobExecuteResult();
result.setStatus(PipelineStatusEnum.SUCCESS);
result.setMessage(PipelineStatusEnum.SUCCESS.getDesc());
result.setData(data);
return result;
}
public static PipelineNodeJobExecuteResult success(String message, Object data) {
PipelineNodeJobExecuteResult result = new PipelineNodeJobExecuteResult();
result.setStatus(PipelineStatusEnum.SUCCESS);
result.setMessage(message);
result.setData(data);
return result;
}
@@ -20,33 +20,33 @@ public class PipelineNodeTemplateVo extends CsObject {
/**
* 业务编码
*/
private String code;
protected String code;
/**
* 业务类型(service:系统内置服务 rpc:远程调用)
*/
private String type;
protected String type;
/**
* 业务名称
*/
private String name;
protected String name;
/**
* 是否可点击
*/
private Boolean click = false;
protected Boolean click = false;
/**
* 是否可重试
*/
private Boolean retry = false;
protected Boolean retry = false;
/**
* 是否可点击:点击展示详情,详情内容类型
*/
private String detailType = "";
protected String detailType = "";
/**
* 默认参数
*/
private Map<String, Object> parameters = new HashMap<>();
protected Map<String, Object> parameters = new HashMap<>();
/**
* 流水线节点控制按钮
*/
private List<PipelineNodeOperateButtonVo> buttons = new ArrayList<>();
protected List<PipelineNodeOperateButtonVo> buttons = new ArrayList<>();
}
@@ -1,34 +0,0 @@
package cn.odboy.devops.job;
import cn.odboy.devops.constant.pipeline.PipelineConst;
import cn.odboy.devops.constant.pipeline.PipelineStatusEnum;
import cn.odboy.devops.framework.pipeline.core.PipelineNodeJobExecutor;
import cn.odboy.devops.framework.pipeline.model.PipelineNodeJobExecuteResult;
import cn.odboy.devops.framework.pipeline.model.PipelineNodeTemplateVo;
import cn.odboy.devops.service.PipelineInstanceNodeDetailService;
import cn.odboy.devops.service.PipelineInstanceNodeService;
import cn.odboy.devops.service.PipelineInstanceService;
import cn.odboy.devops.service.PipelineTemplateService;
import lombok.RequiredArgsConstructor;
import org.quartz.JobDataMap;
import org.springframework.stereotype.Service;
@RequiredArgsConstructor
@Service(value = PipelineConst.EXECUTOR_PREFIX + "node_build_java")
public class PipelineNodeBuildJava implements PipelineNodeJobExecutor {
private final PipelineTemplateService pipelineTemplateService;
private final PipelineInstanceService pipelineInstanceService;
private final PipelineInstanceNodeService pipelineInstanceNodeService;
private final PipelineInstanceNodeDetailService pipelineInstanceNodeDetailService;
@Override
public PipelineNodeJobExecuteResult execute(JobDataMap jobDataMap) {
long instanceId = jobDataMap.getLong(PipelineConst.INSTANCE_ID);
String contextName = jobDataMap.getString(PipelineConst.CONTEXT_NAME);
String env = jobDataMap.getString(PipelineConst.ENV);
long templateId = jobDataMap.getLong(PipelineConst.TEMPLATE_ID);
PipelineNodeTemplateVo currentNodeTemplate = (PipelineNodeTemplateVo) jobDataMap.get(PipelineConst.CURRENT_NODE_TEMPLATE);
pipelineInstanceNodeService.updatePipelineInstanceNodeByArgs(instanceId, currentNodeTemplate.getCode(), PipelineStatusEnum.SUCCESS, PipelineStatusEnum.SUCCESS.getDesc());
return PipelineNodeJobExecuteResult.success();
}
}
@@ -0,0 +1,41 @@
package cn.odboy.devops.job;
import cn.odboy.devops.constant.pipeline.PipelineConst;
import cn.odboy.devops.dal.dataobject.PipelineInstanceNodeTb;
import cn.odboy.devops.framework.pipeline.AbstractPipelineNodeJobService;
import cn.odboy.devops.framework.pipeline.core.PipelineNodeJobExecutor;
import cn.odboy.devops.framework.pipeline.model.PipelineNodeJobExecuteResult;
import cn.odboy.devops.framework.pipeline.model.PipelineNodeTemplateVo;
import cn.odboy.devops.job.biz.PipelineNodeBuildJavaBiz;
import cn.odboy.devops.job.biz.PipelineNodeInitBiz;
import cn.odboy.framework.exception.BadRequestException;
import com.alibaba.fastjson2.JSON;
import lombok.RequiredArgsConstructor;
import org.quartz.JobDataMap;
import org.springframework.stereotype.Service;
import java.util.List;
@RequiredArgsConstructor
@Service(value = PipelineConst.EXECUTOR_PREFIX + "node_build_java")
public class PipelineNodeBuildJavaJob extends AbstractPipelineNodeJobService implements PipelineNodeJobExecutor {
private final PipelineNodeBuildJavaBiz pipelineNodeBuildJavaBiz;
@Override
public PipelineNodeJobExecuteResult execute(JobDataMap jobDataMap) throws BadRequestException {
// 参数列表
long instanceId = jobDataMap.getLong(PipelineConst.INSTANCE_ID);
PipelineNodeTemplateVo currentNodeTemplate = (PipelineNodeTemplateVo) jobDataMap.get(PipelineConst.CURRENT_NODE_TEMPLATE);
PipelineInstanceNodeTb pipelineInstanceNode = getCurrentNodeInfo(instanceId, currentNodeTemplate.getCode());
String contextName = jobDataMap.getString(PipelineConst.CONTEXT_NAME);
String env = jobDataMap.getString(PipelineConst.ENV);
List<PipelineNodeTemplateVo> templateList = JSON.parseArray(jobDataMap.getString(PipelineConst.TEMPLATE), PipelineNodeTemplateVo.class);
PipelineNodeJobExecuteResult lastNodeResult = (PipelineNodeJobExecuteResult) jobDataMap.get(PipelineConst.LAST_NODE_RESULT);
// 步骤执行
pipelineNodeBuildJavaBiz.initStart(pipelineInstanceNode, contextName, env, templateList, lastNodeResult);
pipelineNodeBuildJavaBiz.createReleaseBranch(pipelineInstanceNode, contextName, env, templateList, lastNodeResult);
pipelineNodeBuildJavaBiz.mergeMasterToRelease(pipelineInstanceNode, contextName, env, templateList, lastNodeResult);
pipelineNodeBuildJavaBiz.initFinish(pipelineInstanceNode, contextName, env, templateList, lastNodeResult);
return PipelineNodeJobExecuteResult.success();
}
}
@@ -1,34 +0,0 @@
package cn.odboy.devops.job;
import cn.odboy.devops.constant.pipeline.PipelineConst;
import cn.odboy.devops.constant.pipeline.PipelineStatusEnum;
import cn.odboy.devops.framework.pipeline.core.PipelineNodeJobExecutor;
import cn.odboy.devops.framework.pipeline.model.PipelineNodeJobExecuteResult;
import cn.odboy.devops.framework.pipeline.model.PipelineNodeTemplateVo;
import cn.odboy.devops.service.PipelineInstanceNodeDetailService;
import cn.odboy.devops.service.PipelineInstanceNodeService;
import cn.odboy.devops.service.PipelineInstanceService;
import cn.odboy.devops.service.PipelineTemplateService;
import lombok.RequiredArgsConstructor;
import org.quartz.JobDataMap;
import org.springframework.stereotype.Service;
@RequiredArgsConstructor
@Service(value = PipelineConst.EXECUTOR_PREFIX + "node_deploy_java")
public class PipelineNodeDeployJava implements PipelineNodeJobExecutor {
private final PipelineTemplateService pipelineTemplateService;
private final PipelineInstanceService pipelineInstanceService;
private final PipelineInstanceNodeService pipelineInstanceNodeService;
private final PipelineInstanceNodeDetailService pipelineInstanceNodeDetailService;
@Override
public PipelineNodeJobExecuteResult execute(JobDataMap jobDataMap) {
long instanceId = jobDataMap.getLong(PipelineConst.INSTANCE_ID);
String contextName = jobDataMap.getString(PipelineConst.CONTEXT_NAME);
String env = jobDataMap.getString(PipelineConst.ENV);
long templateId = jobDataMap.getLong(PipelineConst.TEMPLATE_ID);
PipelineNodeTemplateVo currentNodeTemplate = (PipelineNodeTemplateVo) jobDataMap.get(PipelineConst.CURRENT_NODE_TEMPLATE);
pipelineInstanceNodeService.updatePipelineInstanceNodeByArgs(instanceId, currentNodeTemplate.getCode(), PipelineStatusEnum.SUCCESS, PipelineStatusEnum.SUCCESS.getDesc());
return PipelineNodeJobExecuteResult.success();
}
}
@@ -0,0 +1,41 @@
package cn.odboy.devops.job;
import cn.odboy.devops.constant.pipeline.PipelineConst;
import cn.odboy.devops.dal.dataobject.PipelineInstanceNodeTb;
import cn.odboy.devops.framework.pipeline.AbstractPipelineNodeJobService;
import cn.odboy.devops.framework.pipeline.core.PipelineNodeJobExecutor;
import cn.odboy.devops.framework.pipeline.model.PipelineNodeJobExecuteResult;
import cn.odboy.devops.framework.pipeline.model.PipelineNodeTemplateVo;
import cn.odboy.devops.job.biz.PipelineNodeBuildJavaBiz;
import cn.odboy.devops.job.biz.PipelineNodeDeployJavaBiz;
import cn.odboy.framework.exception.BadRequestException;
import com.alibaba.fastjson2.JSON;
import lombok.RequiredArgsConstructor;
import org.quartz.JobDataMap;
import org.springframework.stereotype.Service;
import java.util.List;
@RequiredArgsConstructor
@Service(value = PipelineConst.EXECUTOR_PREFIX + "node_deploy_java")
public class PipelineNodeDeployJavaJob extends AbstractPipelineNodeJobService implements PipelineNodeJobExecutor {
private final PipelineNodeDeployJavaBiz pipelineNodeDeployJavaBiz;
@Override
public PipelineNodeJobExecuteResult execute(JobDataMap jobDataMap) throws BadRequestException {
// 参数列表
long instanceId = jobDataMap.getLong(PipelineConst.INSTANCE_ID);
PipelineNodeTemplateVo currentNodeTemplate = (PipelineNodeTemplateVo) jobDataMap.get(PipelineConst.CURRENT_NODE_TEMPLATE);
PipelineInstanceNodeTb pipelineInstanceNode = getCurrentNodeInfo(instanceId, currentNodeTemplate.getCode());
String contextName = jobDataMap.getString(PipelineConst.CONTEXT_NAME);
String env = jobDataMap.getString(PipelineConst.ENV);
List<PipelineNodeTemplateVo> templateList = JSON.parseArray(jobDataMap.getString(PipelineConst.TEMPLATE), PipelineNodeTemplateVo.class);
PipelineNodeJobExecuteResult lastNodeResult = (PipelineNodeJobExecuteResult) jobDataMap.get(PipelineConst.LAST_NODE_RESULT);
// 步骤执行
pipelineNodeDeployJavaBiz.initStart(pipelineInstanceNode, contextName, env, templateList, lastNodeResult);
pipelineNodeDeployJavaBiz.createReleaseBranch(pipelineInstanceNode, contextName, env, templateList, lastNodeResult);
pipelineNodeDeployJavaBiz.mergeMasterToRelease(pipelineInstanceNode, contextName, env, templateList, lastNodeResult);
pipelineNodeDeployJavaBiz.initFinish(pipelineInstanceNode, contextName, env, templateList, lastNodeResult);
return PipelineNodeJobExecuteResult.success();
}
}
@@ -1,34 +1,42 @@
package cn.odboy.devops.job;
import cn.odboy.devops.constant.pipeline.PipelineConst;
import cn.odboy.devops.constant.pipeline.PipelineStatusEnum;
import cn.odboy.devops.dal.dataobject.PipelineInstanceNodeTb;
import cn.odboy.devops.framework.pipeline.AbstractPipelineNodeJobService;
import cn.odboy.devops.framework.pipeline.core.PipelineNodeJobExecutor;
import cn.odboy.devops.framework.pipeline.model.PipelineNodeJobExecuteResult;
import cn.odboy.devops.framework.pipeline.model.PipelineNodeTemplateVo;
import cn.odboy.devops.service.PipelineInstanceNodeDetailService;
import cn.odboy.devops.service.PipelineInstanceNodeService;
import cn.odboy.devops.service.PipelineInstanceService;
import cn.odboy.devops.service.PipelineTemplateService;
import cn.odboy.devops.job.biz.PipelineNodeInitBiz;
import cn.odboy.framework.exception.BadRequestException;
import com.alibaba.fastjson2.JSON;
import lombok.RequiredArgsConstructor;
import lombok.extern.slf4j.Slf4j;
import org.quartz.JobDataMap;
import org.springframework.stereotype.Service;
@RequiredArgsConstructor
import java.util.List;
@Slf4j
@Service(value = PipelineConst.EXECUTOR_PREFIX + "node_init")
public class PipelineNodeInitJob implements PipelineNodeJobExecutor {
private final PipelineTemplateService pipelineTemplateService;
private final PipelineInstanceService pipelineInstanceService;
private final PipelineInstanceNodeService pipelineInstanceNodeService;
private final PipelineInstanceNodeDetailService pipelineInstanceNodeDetailService;
@RequiredArgsConstructor
public class PipelineNodeInitJob extends AbstractPipelineNodeJobService implements PipelineNodeJobExecutor {
private final PipelineNodeInitBiz pipelineNodeInitBiz;
@Override
public PipelineNodeJobExecuteResult execute(JobDataMap jobDataMap) {
public PipelineNodeJobExecuteResult execute(JobDataMap jobDataMap) throws BadRequestException {
// 参数列表
long instanceId = jobDataMap.getLong(PipelineConst.INSTANCE_ID);
PipelineNodeTemplateVo currentNodeTemplate = (PipelineNodeTemplateVo) jobDataMap.get(PipelineConst.CURRENT_NODE_TEMPLATE);
PipelineInstanceNodeTb pipelineInstanceNode = getCurrentNodeInfo(instanceId, currentNodeTemplate.getCode());
String contextName = jobDataMap.getString(PipelineConst.CONTEXT_NAME);
String env = jobDataMap.getString(PipelineConst.ENV);
long templateId = jobDataMap.getLong(PipelineConst.TEMPLATE_ID);
PipelineNodeTemplateVo currentNodeTemplate = (PipelineNodeTemplateVo) jobDataMap.get(PipelineConst.CURRENT_NODE_TEMPLATE);
pipelineInstanceNodeService.updatePipelineInstanceNodeByArgs(instanceId, currentNodeTemplate.getCode(), PipelineStatusEnum.SUCCESS, PipelineStatusEnum.SUCCESS.getDesc());
List<PipelineNodeTemplateVo> templateList = JSON.parseArray(jobDataMap.getString(PipelineConst.TEMPLATE), PipelineNodeTemplateVo.class);
PipelineNodeJobExecuteResult lastNodeResult = (PipelineNodeJobExecuteResult) jobDataMap.get(PipelineConst.LAST_NODE_RESULT);
// 步骤执行
pipelineNodeInitBiz.initStart(pipelineInstanceNode, contextName, env, templateList, lastNodeResult);
pipelineNodeInitBiz.createReleaseBranch(pipelineInstanceNode, contextName, env, templateList, lastNodeResult);
pipelineNodeInitBiz.mergeMasterToRelease(pipelineInstanceNode, contextName, env, templateList, lastNodeResult);
pipelineNodeInitBiz.initFinish(pipelineInstanceNode, contextName, env, templateList, lastNodeResult);
return PipelineNodeJobExecuteResult.success();
}
}
@@ -1,34 +0,0 @@
package cn.odboy.devops.job;
import cn.odboy.devops.constant.pipeline.PipelineConst;
import cn.odboy.devops.constant.pipeline.PipelineStatusEnum;
import cn.odboy.devops.framework.pipeline.core.PipelineNodeJobExecutor;
import cn.odboy.devops.framework.pipeline.model.PipelineNodeJobExecuteResult;
import cn.odboy.devops.framework.pipeline.model.PipelineNodeTemplateVo;
import cn.odboy.devops.service.PipelineInstanceNodeDetailService;
import cn.odboy.devops.service.PipelineInstanceNodeService;
import cn.odboy.devops.service.PipelineInstanceService;
import cn.odboy.devops.service.PipelineTemplateService;
import lombok.RequiredArgsConstructor;
import org.quartz.JobDataMap;
import org.springframework.stereotype.Service;
@RequiredArgsConstructor
@Service(value = PipelineConst.EXECUTOR_PREFIX + "node_merge_branch")
public class PipelineNodeMergeBranch implements PipelineNodeJobExecutor {
private final PipelineTemplateService pipelineTemplateService;
private final PipelineInstanceService pipelineInstanceService;
private final PipelineInstanceNodeService pipelineInstanceNodeService;
private final PipelineInstanceNodeDetailService pipelineInstanceNodeDetailService;
@Override
public PipelineNodeJobExecuteResult execute(JobDataMap jobDataMap) {
long instanceId = jobDataMap.getLong(PipelineConst.INSTANCE_ID);
String contextName = jobDataMap.getString(PipelineConst.CONTEXT_NAME);
String env = jobDataMap.getString(PipelineConst.ENV);
long templateId = jobDataMap.getLong(PipelineConst.TEMPLATE_ID);
PipelineNodeTemplateVo currentNodeTemplate = (PipelineNodeTemplateVo) jobDataMap.get(PipelineConst.CURRENT_NODE_TEMPLATE);
pipelineInstanceNodeService.updatePipelineInstanceNodeByArgs(instanceId, currentNodeTemplate.getCode(), PipelineStatusEnum.SUCCESS, PipelineStatusEnum.SUCCESS.getDesc());
return PipelineNodeJobExecuteResult.success();
}
}
@@ -0,0 +1,41 @@
package cn.odboy.devops.job;
import cn.odboy.devops.constant.pipeline.PipelineConst;
import cn.odboy.devops.dal.dataobject.PipelineInstanceNodeTb;
import cn.odboy.devops.framework.pipeline.AbstractPipelineNodeJobService;
import cn.odboy.devops.framework.pipeline.core.PipelineNodeJobExecutor;
import cn.odboy.devops.framework.pipeline.model.PipelineNodeJobExecuteResult;
import cn.odboy.devops.framework.pipeline.model.PipelineNodeTemplateVo;
import cn.odboy.devops.job.biz.PipelineNodeInitBiz;
import cn.odboy.devops.job.biz.PipelineNodeMergeBranchBiz;
import cn.odboy.framework.exception.BadRequestException;
import com.alibaba.fastjson2.JSON;
import lombok.RequiredArgsConstructor;
import org.quartz.JobDataMap;
import org.springframework.stereotype.Service;
import java.util.List;
@RequiredArgsConstructor
@Service(value = PipelineConst.EXECUTOR_PREFIX + "node_merge_branch")
public class PipelineNodeMergeBranchJob extends AbstractPipelineNodeJobService implements PipelineNodeJobExecutor {
private final PipelineNodeMergeBranchBiz pipelineNodeMergeBranchBiz;
@Override
public PipelineNodeJobExecuteResult execute(JobDataMap jobDataMap) throws BadRequestException {
// 参数列表
long instanceId = jobDataMap.getLong(PipelineConst.INSTANCE_ID);
PipelineNodeTemplateVo currentNodeTemplate = (PipelineNodeTemplateVo) jobDataMap.get(PipelineConst.CURRENT_NODE_TEMPLATE);
PipelineInstanceNodeTb pipelineInstanceNode = getCurrentNodeInfo(instanceId, currentNodeTemplate.getCode());
String contextName = jobDataMap.getString(PipelineConst.CONTEXT_NAME);
String env = jobDataMap.getString(PipelineConst.ENV);
List<PipelineNodeTemplateVo> templateList = JSON.parseArray(jobDataMap.getString(PipelineConst.TEMPLATE), PipelineNodeTemplateVo.class);
PipelineNodeJobExecuteResult lastNodeResult = (PipelineNodeJobExecuteResult) jobDataMap.get(PipelineConst.LAST_NODE_RESULT);
// 步骤执行
pipelineNodeMergeBranchBiz.initStart(pipelineInstanceNode, contextName, env, templateList, lastNodeResult);
pipelineNodeMergeBranchBiz.createReleaseBranch(pipelineInstanceNode, contextName, env, templateList, lastNodeResult);
pipelineNodeMergeBranchBiz.mergeMasterToRelease(pipelineInstanceNode, contextName, env, templateList, lastNodeResult);
pipelineNodeMergeBranchBiz.initFinish(pipelineInstanceNode, contextName, env, templateList, lastNodeResult);
return PipelineNodeJobExecuteResult.success();
}
}
@@ -0,0 +1,35 @@
package cn.odboy.devops.job.biz;
import cn.hutool.core.thread.ThreadUtil;
import cn.odboy.devops.dal.dataobject.PipelineInstanceNodeTb;
import cn.odboy.devops.framework.pipeline.log.PipelineNodeStepLog;
import cn.odboy.devops.framework.pipeline.model.PipelineNodeJobExecuteResult;
import cn.odboy.devops.framework.pipeline.model.PipelineNodeTemplateVo;
import lombok.extern.slf4j.Slf4j;
import org.springframework.stereotype.Service;
import java.util.List;
@Slf4j
@Service
public class PipelineNodeBuildJavaBiz {
@PipelineNodeStepLog("初始化开始")
public void initStart(PipelineInstanceNodeTb pipelineInstanceNode, String contextName, String env, List<PipelineNodeTemplateVo> templateList, PipelineNodeJobExecuteResult lastNodeResult) {
ThreadUtil.safeSleep(5000);
}
@PipelineNodeStepLog("Master分支合并到Release分支")
public void mergeMasterToRelease(PipelineInstanceNodeTb pipelineInstanceNode, String contextName, String env, List<PipelineNodeTemplateVo> templateList, PipelineNodeJobExecuteResult lastNodeResult) {
ThreadUtil.safeSleep(5000);
}
@PipelineNodeStepLog("新建Release分支")
public void createReleaseBranch(PipelineInstanceNodeTb pipelineInstanceNode, String contextName, String env, List<PipelineNodeTemplateVo> templateList, PipelineNodeJobExecuteResult lastNodeResult) {
ThreadUtil.safeSleep(5000);
}
@PipelineNodeStepLog("初始化完成")
public void initFinish(PipelineInstanceNodeTb pipelineInstanceNode, String contextName, String env, List<PipelineNodeTemplateVo> templateList, PipelineNodeJobExecuteResult lastNodeResult) {
ThreadUtil.safeSleep(5000);
}
}
@@ -0,0 +1,35 @@
package cn.odboy.devops.job.biz;
import cn.hutool.core.thread.ThreadUtil;
import cn.odboy.devops.dal.dataobject.PipelineInstanceNodeTb;
import cn.odboy.devops.framework.pipeline.log.PipelineNodeStepLog;
import cn.odboy.devops.framework.pipeline.model.PipelineNodeJobExecuteResult;
import cn.odboy.devops.framework.pipeline.model.PipelineNodeTemplateVo;
import lombok.extern.slf4j.Slf4j;
import org.springframework.stereotype.Service;
import java.util.List;
@Slf4j
@Service
public class PipelineNodeDeployJavaBiz {
@PipelineNodeStepLog("初始化开始")
public void initStart(PipelineInstanceNodeTb pipelineInstanceNode, String contextName, String env, List<PipelineNodeTemplateVo> templateList, PipelineNodeJobExecuteResult lastNodeResult) {
ThreadUtil.safeSleep(5000);
}
@PipelineNodeStepLog("Master分支合并到Release分支")
public void mergeMasterToRelease(PipelineInstanceNodeTb pipelineInstanceNode, String contextName, String env, List<PipelineNodeTemplateVo> templateList, PipelineNodeJobExecuteResult lastNodeResult) {
ThreadUtil.safeSleep(5000);
}
@PipelineNodeStepLog("新建Release分支")
public void createReleaseBranch(PipelineInstanceNodeTb pipelineInstanceNode, String contextName, String env, List<PipelineNodeTemplateVo> templateList, PipelineNodeJobExecuteResult lastNodeResult) {
ThreadUtil.safeSleep(5000);
}
@PipelineNodeStepLog("初始化完成")
public void initFinish(PipelineInstanceNodeTb pipelineInstanceNode, String contextName, String env, List<PipelineNodeTemplateVo> templateList, PipelineNodeJobExecuteResult lastNodeResult) {
ThreadUtil.safeSleep(5000);
}
}
@@ -0,0 +1,35 @@
package cn.odboy.devops.job.biz;
import cn.hutool.core.thread.ThreadUtil;
import cn.odboy.devops.dal.dataobject.PipelineInstanceNodeTb;
import cn.odboy.devops.framework.pipeline.log.PipelineNodeStepLog;
import cn.odboy.devops.framework.pipeline.model.PipelineNodeJobExecuteResult;
import cn.odboy.devops.framework.pipeline.model.PipelineNodeTemplateVo;
import lombok.extern.slf4j.Slf4j;
import org.springframework.stereotype.Service;
import java.util.List;
@Slf4j
@Service
public class PipelineNodeInitBiz {
@PipelineNodeStepLog("初始化开始")
public void initStart(PipelineInstanceNodeTb pipelineInstanceNode, String contextName, String env, List<PipelineNodeTemplateVo> templateList, PipelineNodeJobExecuteResult lastNodeResult) {
ThreadUtil.safeSleep(5000);
}
@PipelineNodeStepLog("Master分支合并到Release分支")
public void mergeMasterToRelease(PipelineInstanceNodeTb pipelineInstanceNode, String contextName, String env, List<PipelineNodeTemplateVo> templateList, PipelineNodeJobExecuteResult lastNodeResult) {
ThreadUtil.safeSleep(5000);
}
@PipelineNodeStepLog("新建Release分支")
public void createReleaseBranch(PipelineInstanceNodeTb pipelineInstanceNode, String contextName, String env, List<PipelineNodeTemplateVo> templateList, PipelineNodeJobExecuteResult lastNodeResult) {
ThreadUtil.safeSleep(5000);
}
@PipelineNodeStepLog("初始化完成")
public void initFinish(PipelineInstanceNodeTb pipelineInstanceNode, String contextName, String env, List<PipelineNodeTemplateVo> templateList, PipelineNodeJobExecuteResult lastNodeResult) {
ThreadUtil.safeSleep(5000);
}
}
@@ -0,0 +1,35 @@
package cn.odboy.devops.job.biz;
import cn.hutool.core.thread.ThreadUtil;
import cn.odboy.devops.dal.dataobject.PipelineInstanceNodeTb;
import cn.odboy.devops.framework.pipeline.log.PipelineNodeStepLog;
import cn.odboy.devops.framework.pipeline.model.PipelineNodeJobExecuteResult;
import cn.odboy.devops.framework.pipeline.model.PipelineNodeTemplateVo;
import lombok.extern.slf4j.Slf4j;
import org.springframework.stereotype.Service;
import java.util.List;
@Slf4j
@Service
public class PipelineNodeMergeBranchBiz {
@PipelineNodeStepLog("初始化开始")
public void initStart(PipelineInstanceNodeTb pipelineInstanceNode, String contextName, String env, List<PipelineNodeTemplateVo> templateList, PipelineNodeJobExecuteResult lastNodeResult) {
ThreadUtil.safeSleep(5000);
}
@PipelineNodeStepLog("Master分支合并到Release分支")
public void mergeMasterToRelease(PipelineInstanceNodeTb pipelineInstanceNode, String contextName, String env, List<PipelineNodeTemplateVo> templateList, PipelineNodeJobExecuteResult lastNodeResult) {
ThreadUtil.safeSleep(5000);
}
@PipelineNodeStepLog("新建Release分支")
public void createReleaseBranch(PipelineInstanceNodeTb pipelineInstanceNode, String contextName, String env, List<PipelineNodeTemplateVo> templateList, PipelineNodeJobExecuteResult lastNodeResult) {
ThreadUtil.safeSleep(5000);
}
@PipelineNodeStepLog("初始化完成")
public void initFinish(PipelineInstanceNodeTb pipelineInstanceNode, String contextName, String env, List<PipelineNodeTemplateVo> templateList, PipelineNodeJobExecuteResult lastNodeResult) {
ThreadUtil.safeSleep(5000);
}
}
@@ -1,8 +1,13 @@
package cn.odboy.devops.service;
import cn.odboy.devops.constant.pipeline.PipelineStatusEnum;
import cn.odboy.devops.dal.dataobject.PipelineInstanceNodeDetailTb;
import cn.odboy.devops.dal.dataobject.PipelineInstanceNodeTb;
import com.anwen.mongo.service.IService;
import java.util.Date;
import java.util.List;
/**
* <p>
* 流水线实例节点明细
@@ -12,5 +17,11 @@ import com.anwen.mongo.service.IService;
* @since 2025-06-26
*/
public interface PipelineInstanceNodeDetailService extends IService<PipelineInstanceNodeDetailTb> {
void addLog(Long nodeId, String nodeCode, String stepName, PipelineStatusEnum status, String stepMsg, Date finishTime);
PipelineInstanceNodeDetailTb getPipelineInstanceNodeDetailByArgs(Long nodeId, String nodeCode, String stepName);
void removeByNodeIds(List<Long> nodeIds);
PipelineInstanceNodeDetailTb getLastPipelineInstanceNodeDetailByArgs(PipelineInstanceNodeTb instanceNode, Long nodeId, String nodeCode);
}
@@ -5,6 +5,9 @@ import cn.odboy.devops.dal.dataobject.PipelineInstanceNodeTb;
import cn.odboy.devops.dal.dataobject.PipelineInstanceTb;
import com.anwen.mongo.service.IService;
import java.util.Date;
import java.util.List;
/**
* 流水线实例节点明细
*
@@ -13,6 +16,14 @@ import com.anwen.mongo.service.IService;
*/
public interface PipelineInstanceNodeService extends IService<PipelineInstanceNodeTb> {
void createPipelineInstanceDetail(PipelineInstanceTb pipelineInstanceTb);
void updatePipelineInstanceNodeByArgs(Long instanceId, String code, PipelineStatusEnum pipelineStatusEnum, String msg);
PipelineInstanceNodeTb getPipelineInstanceNodeByArgs(Long instanceId, String code);
void remakePipelineInstanceDetail(PipelineInstanceTb pipelineInstanceTb, String retryNodeCode);
List<PipelineInstanceNodeTb> queryPipelineInstanceNodeListByInstanceId(Long instanceId);
void finishPipelineInstanceNodeByArgs(Long instanceId, String nodeCode, PipelineStatusEnum pipelineStatusEnum, String desc, Date date);
}
@@ -1,6 +1,9 @@
package cn.odboy.devops.service;
import cn.odboy.devops.dal.dataobject.PipelineInstanceTb;
import cn.odboy.devops.framework.pipeline.model.PipelineNodeDataVo;
import java.util.List;
/**
* <p>
@@ -11,5 +14,7 @@ import cn.odboy.devops.dal.dataobject.PipelineInstanceTb;
* @since 2025-06-26
*/
public interface PipelineInstanceService {
void start(PipelineInstanceTb pipelineInstanceTb);
String startPipeline(PipelineInstanceTb pipelineInstanceTb);
void restartPipeline(PipelineInstanceTb pipelineInstanceTb, String retryNodeCode);
List<PipelineNodeDataVo> queryLastPipelineDetail(String instanceId);
}
@@ -1,31 +1,36 @@
package cn.odboy.devops.service.impl;
import cn.hutool.core.collection.CollUtil;
import cn.odboy.devops.constant.pipeline.PipelineStatusEnum;
import cn.odboy.devops.dal.dataobject.PipelineInstanceNodeDetailTb;
import cn.odboy.devops.dal.dataobject.PipelineInstanceNodeTb;
import cn.odboy.devops.service.PipelineInstanceNodeDetailService;
import com.anwen.mongo.conditions.query.LambdaQueryChainWrapper;
import com.anwen.mongo.conditions.update.LambdaUpdateChainWrapper;
import com.anwen.mongo.service.impl.ServiceImpl;
import lombok.extern.slf4j.Slf4j;
import org.springframework.stereotype.Service;
import java.util.Date;
import java.util.List;
@Slf4j
@Service
public class PipelineInstanceNodeDetailServiceImpl extends ServiceImpl<PipelineInstanceNodeDetailTb> implements PipelineInstanceNodeDetailService {
public void addLog(String nodeId, Integer nodeIndex, String nodeCode, String stepName, PipelineStatusEnum status, String stepMsg) {
@Override
public void addLog(Long nodeId, String nodeCode, String stepName, PipelineStatusEnum status, String stepMsg, Date finishTime) {
PipelineInstanceNodeDetailTb record = getPipelineInstanceNodeDetailByArgs(nodeId, nodeCode, stepName);
if (record != null) {
record.setStepStatus(status.getCode());
record.setStepMsg(stepMsg);
record.setFinishTime(finishTime);
updateById(record);
return;
}
record = new PipelineInstanceNodeDetailTb();
record.setNodeId(nodeId);
record.setStartTime(new Date());
record.setNodeIndex(nodeIndex);
record.setNodeCode(nodeCode);
record.setStepName(stepName);
record.setStepStatus(status.getCode());
@@ -33,10 +38,53 @@ public class PipelineInstanceNodeDetailServiceImpl extends ServiceImpl<PipelineI
save(record);
}
public PipelineInstanceNodeDetailTb getPipelineInstanceNodeDetailByArgs(String nodeId, String nodeCode, String stepName) {
return one(new LambdaQueryChainWrapper<>(getBaseMapper(), PipelineInstanceNodeDetailTb.class)
.eq(PipelineInstanceNodeDetailTb::getNodeId, nodeId).eq(PipelineInstanceNodeDetailTb::getNodeCode, nodeCode)
.eq(PipelineInstanceNodeDetailTb::getStepName, stepName)
@Override
public PipelineInstanceNodeDetailTb getPipelineInstanceNodeDetailByArgs(Long nodeId, String nodeCode, String stepName) {
return one(new LambdaQueryChainWrapper<>(getBaseMapper(), PipelineInstanceNodeDetailTb.class).eq(PipelineInstanceNodeDetailTb::getNodeId, nodeId).eq(PipelineInstanceNodeDetailTb::getNodeCode, nodeCode).eq(PipelineInstanceNodeDetailTb::getStepName, stepName));
}
@Override
public void removeByNodeIds(List<Long> nodeIds) {
if (CollUtil.isNotEmpty(nodeIds)) {
remove(new LambdaUpdateChainWrapper<>(getBaseMapper(), PipelineInstanceNodeDetailTb.class).in(PipelineInstanceNodeDetailTb::getNodeId, nodeIds));
}
}
@Override
public PipelineInstanceNodeDetailTb getLastPipelineInstanceNodeDetailByArgs(PipelineInstanceNodeTb instanceNode, Long nodeId, String nodeCode) {
List<PipelineInstanceNodeDetailTb> pipelineInstanceNodeDetailList = list(new LambdaQueryChainWrapper<>(getBaseMapper(), PipelineInstanceNodeDetailTb.class)
.eq(PipelineInstanceNodeDetailTb::getNodeId, nodeId)
.eq(PipelineInstanceNodeDetailTb::getNodeCode, nodeCode)
.orderByDesc(PipelineInstanceNodeDetailTb::getId)
);
// 明细为空,返回节点状态
if (pipelineInstanceNodeDetailList.isEmpty()) {
PipelineInstanceNodeDetailTb pipelineInstanceNodeDetailTb = new PipelineInstanceNodeDetailTb();
pipelineInstanceNodeDetailTb.setStepMsg(instanceNode.getCurrentNodeMsg());
pipelineInstanceNodeDetailTb.setStepStatus(instanceNode.getCurrentNodeStatus());
return pipelineInstanceNodeDetailTb;
}
boolean existFail = pipelineInstanceNodeDetailList.stream().anyMatch(f -> PipelineStatusEnum.FAIL.getCode().equals(f.getStepStatus()));
// 节点中存在失败的步骤
if (existFail) {
PipelineInstanceNodeDetailTb pipelineInstanceNodeDetailTb = new PipelineInstanceNodeDetailTb();
pipelineInstanceNodeDetailTb.setStepMsg(instanceNode.getCurrentNodeMsg());
pipelineInstanceNodeDetailTb.setStepStatus(instanceNode.getCurrentNodeStatus());
return pipelineInstanceNodeDetailTb;
}
// 节点中所有步骤都成功
long totalStep = pipelineInstanceNodeDetailList.size();
long totalSuccessStep = pipelineInstanceNodeDetailList.stream().filter(f -> PipelineStatusEnum.SUCCESS.getCode().equals(f.getStepStatus())).count();
if (totalSuccessStep == totalStep) {
PipelineInstanceNodeDetailTb pipelineInstanceNodeDetailTb = new PipelineInstanceNodeDetailTb();
pipelineInstanceNodeDetailTb.setStepMsg(instanceNode.getCurrentNodeMsg());
pipelineInstanceNodeDetailTb.setStepStatus(instanceNode.getCurrentNodeStatus());
return pipelineInstanceNodeDetailTb;
}
// 返回节点步骤状态
PipelineInstanceNodeDetailTb pipelineInstanceNodeDetailTb = pipelineInstanceNodeDetailList.get(0);
// 返回步骤说明
pipelineInstanceNodeDetailTb.setStepMsg(pipelineInstanceNodeDetailTb.getStepName());
return pipelineInstanceNodeDetailTb;
}
}
@@ -5,20 +5,26 @@ import cn.odboy.devops.constant.pipeline.PipelineStatusEnum;
import cn.odboy.devops.dal.dataobject.PipelineInstanceNodeTb;
import cn.odboy.devops.dal.dataobject.PipelineInstanceTb;
import cn.odboy.devops.framework.pipeline.model.PipelineNodeTemplateVo;
import cn.odboy.devops.service.PipelineInstanceNodeDetailService;
import cn.odboy.devops.service.PipelineInstanceNodeService;
import com.alibaba.fastjson2.JSON;
import com.anwen.mongo.conditions.query.LambdaQueryChainWrapper;
import com.anwen.mongo.service.impl.ServiceImpl;
import lombok.RequiredArgsConstructor;
import lombok.extern.slf4j.Slf4j;
import org.springframework.stereotype.Service;
import java.util.ArrayList;
import java.util.Date;
import java.util.List;
import java.util.stream.Collectors;
@Slf4j
@Service
@RequiredArgsConstructor
public class PipelineInstanceNodeServiceImpl extends ServiceImpl<PipelineInstanceNodeTb> implements PipelineInstanceNodeService {
private final PipelineInstanceNodeDetailService pipelineInstanceNodeDetailService;
@Override
public void createPipelineInstanceDetail(PipelineInstanceTb pipelineInstanceTb) {
String templateContent = pipelineInstanceTb.getTemplateContent();
@@ -36,6 +42,8 @@ public class PipelineInstanceNodeServiceImpl extends ServiceImpl<PipelineInstanc
record.setClick(pipelineNodeTemplateVo.getClick());
record.setRetry(pipelineNodeTemplateVo.getRetry());
record.setButtons(pipelineNodeTemplateVo.getButtons());
record.setCurrentNodeStatus(PipelineStatusEnum.PENDING.getCode());
record.setCurrentNodeMsg(PipelineStatusEnum.PENDING.getDesc());
records.add(record);
}
saveBatch(records);
@@ -57,4 +65,55 @@ public class PipelineInstanceNodeServiceImpl extends ServiceImpl<PipelineInstanc
.eq(PipelineInstanceNodeTb::getCode, code)
);
}
@Override
public void remakePipelineInstanceDetail(PipelineInstanceTb pipelineInstanceTb, String retryNodeCode) {
List<String> penddingNodeList = new ArrayList<>();
List<PipelineNodeTemplateVo> pipelineNodeTemplateVos = JSON.parseArray(pipelineInstanceTb.getTemplateContent(), PipelineNodeTemplateVo.class);
boolean findFlag = false;
for (PipelineNodeTemplateVo pipelineNodeTemplateVo : pipelineNodeTemplateVos) {
if (pipelineNodeTemplateVo.getCode().equals(retryNodeCode) || findFlag) {
findFlag = true;
penddingNodeList.add(pipelineNodeTemplateVo.getCode());
}
}
// 更新节点状态
List<PipelineInstanceNodeTb> instanceNodeList = list(new LambdaQueryChainWrapper<>(getBaseMapper(), PipelineInstanceNodeTb.class)
.eq(PipelineInstanceNodeTb::getInstanceId, pipelineInstanceTb.getInstanceId())
.eq(PipelineInstanceNodeTb::getCode, penddingNodeList)
);
if (CollUtil.isNotEmpty(instanceNodeList)) {
Date nowTime = new Date();
for (PipelineInstanceNodeTb pipelineInstanceNodeTb : instanceNodeList) {
pipelineInstanceNodeTb.setUpdateTime(nowTime);
pipelineInstanceNodeTb.setStartTime(null);
pipelineInstanceNodeTb.setCurrentNodeMsg(PipelineStatusEnum.PENDING.getDesc());
pipelineInstanceNodeTb.setCurrentNodeStatus(PipelineStatusEnum.PENDING.getCode());
}
updateBatchByIds(instanceNodeList);
// 删除节点明细数据
List<Long> nodeIds = instanceNodeList.stream()
.map(PipelineInstanceNodeTb::getId)
.distinct()
.collect(Collectors.toList());
pipelineInstanceNodeDetailService.removeByNodeIds(nodeIds);
}
}
@Override
public List<PipelineInstanceNodeTb> queryPipelineInstanceNodeListByInstanceId(Long instanceId) {
return list(new LambdaQueryChainWrapper<>(getBaseMapper(), PipelineInstanceNodeTb.class)
.eq(PipelineInstanceNodeTb::getInstanceId, instanceId)
.orderByAsc(PipelineInstanceNodeTb::getId)
);
}
@Override
public void finishPipelineInstanceNodeByArgs(Long instanceId, String nodeCode, PipelineStatusEnum pipelineStatusEnum, String desc, Date date) {
PipelineInstanceNodeTb pipelineInstanceNode = getPipelineInstanceNodeByArgs(instanceId, nodeCode);
pipelineInstanceNode.setCurrentNodeStatus(pipelineStatusEnum.getCode());
pipelineInstanceNode.setCurrentNodeMsg(desc);
pipelineInstanceNode.setFinishTime(date);
updateById(pipelineInstanceNode);
}
}
@@ -1,14 +1,15 @@
package cn.odboy.devops.service.impl;
import cn.hutool.core.thread.ThreadUtil;
import cn.hutool.core.util.IdUtil;
import cn.hutool.core.bean.BeanUtil;
import cn.odboy.devops.constant.pipeline.PipelineConst;
import cn.odboy.devops.constant.pipeline.PipelineStatusEnum;
import cn.odboy.devops.dal.dataobject.PipelineInstanceNodeDetailTb;
import cn.odboy.devops.dal.dataobject.PipelineInstanceNodeTb;
import cn.odboy.devops.dal.dataobject.PipelineInstanceTb;
import cn.odboy.devops.dal.mysql.PipelineInstanceMapper;
import cn.odboy.devops.dal.redis.PipelineInstanceDAO;
import cn.odboy.devops.framework.pipeline.core.PipelineJobManage;
import cn.odboy.devops.framework.pipeline.model.PipelineNodeTemplateVo;
import cn.odboy.devops.framework.pipeline.model.PipelineNodeDataVo;
import cn.odboy.devops.service.PipelineInstanceNodeDetailService;
import cn.odboy.devops.service.PipelineInstanceNodeService;
import cn.odboy.devops.service.PipelineInstanceService;
@@ -18,6 +19,7 @@ import lombok.RequiredArgsConstructor;
import lombok.extern.slf4j.Slf4j;
import org.springframework.stereotype.Service;
import java.util.ArrayList;
import java.util.List;
@Slf4j
@@ -26,89 +28,63 @@ import java.util.List;
public class PipelineInstanceServiceImpl implements PipelineInstanceService {
private final PipelineJobManage pipelineJobManage;
private final PipelineInstanceMapper pipelineInstanceMapper;
private final PipelineInstanceDAO pipelineInstanceDAO;
private final PipelineInstanceNodeService pipelineInstanceNodeService;
private final PipelineInstanceNodeDetailService pipelineInstanceNodeDetailService;
private final PipelineInstanceDAO pipelineInstanceDAO;
@Override
public void start(PipelineInstanceTb pipelineInstanceTb) {
public String startPipeline(PipelineInstanceTb pipelineInstanceTb) {
if (pipelineInstanceDAO.lock(pipelineInstanceTb)) {
throw new BadRequestException("流水线运行中,无法重复执行");
}
ThreadUtil.execAsync(() -> innerStart(pipelineInstanceTb));
return PipelineConst.INSTANCE_ID + pipelineJobManage.startJob(pipelineInstanceTb);
}
public void innerStart(PipelineInstanceTb pipelineInstanceTb) {
// 雪花id
pipelineInstanceTb.setInstanceId(IdUtil.getSnowflakeNextId());
pipelineInstanceTb.setStatus(PipelineStatusEnum.PENDING.getCode());
// 创建实例
pipelineInstanceMapper.insert(pipelineInstanceTb);
// 创建实例明细
pipelineInstanceNodeService.createPipelineInstanceDetail(pipelineInstanceTb);
// 流水线节点模板
String templateContent = pipelineInstanceTb.getTemplateContent();
List<PipelineNodeTemplateVo> pipelineNodeTemplateVos = JSON.parseArray(templateContent, PipelineNodeTemplateVo.class);
Long instanceId = pipelineInstanceTb.getInstanceId();
// 流水线实例状态(默认是成功的)
PipelineStatusEnum pipelineStatusEnum = PipelineStatusEnum.SUCCESS;
try {
pipelineInstanceTb.setStatus(PipelineStatusEnum.RUNNING.getCode());
pipelineInstanceMapper.updateById(pipelineInstanceTb);
// 依次执行流水线节点任务
for (PipelineNodeTemplateVo pipelineNodeTemplateVo : pipelineNodeTemplateVos) {
try {
pipelineJobManage.startJob(pipelineInstanceTb, pipelineNodeTemplateVo);
pipelineInstanceTb.setCurrentNode(pipelineNodeTemplateVo.getCode());
pipelineInstanceTb.setCurrentNodeStatus(PipelineStatusEnum.RUNNING.getCode());
pipelineInstanceMapper.updateById(pipelineInstanceTb);
// 监听节点执行状态
while (true) {
ThreadUtil.safeSleep(2000);
PipelineInstanceNodeTb pipelineInstanceNode = pipelineInstanceNodeService.getPipelineInstanceNodeByArgs(instanceId, pipelineNodeTemplateVo.getCode());
if (PipelineStatusEnum.SUCCESS.getCode().equals(pipelineInstanceNode.getCurrentNodeStatus())) {
pipelineStatusEnum = PipelineStatusEnum.SUCCESS;
pipelineInstanceTb.setCurrentNodeStatus(pipelineStatusEnum.getCode());
pipelineInstanceMapper.updateById(pipelineInstanceTb);
break;
@Override
public void restartPipeline(PipelineInstanceTb pipelineInstanceTb, String retryNodeCode) {
if (pipelineInstanceTb.getInstanceId() == null) {
throw new BadRequestException("流水线实例Id必填");
}
if (PipelineStatusEnum.FAIL.getCode().equals(pipelineInstanceNode.getCurrentNodeStatus())) {
pipelineStatusEnum = PipelineStatusEnum.FAIL;
pipelineInstanceTb.setCurrentNodeStatus(pipelineStatusEnum.getCode());
pipelineInstanceMapper.updateById(pipelineInstanceTb);
break;
if (pipelineInstanceDAO.lock(pipelineInstanceTb)) {
throw new BadRequestException("流水线运行中,无法重复执行");
}
PipelineInstanceTb currentInstance = pipelineInstanceMapper.selectById(pipelineInstanceTb.getInstanceId());
if (currentInstance == null) {
throw new BadRequestException("无效流水线,请刷新页面后再试");
}
if (pipelineStatusEnum.equals(PipelineStatusEnum.FAIL)) {
// 节点失败,全失败
break;
}
ThreadUtil.safeSleep(2000);
} catch (Exception e) {
log.error("流水线节点运行失败", e);
pipelineInstanceTb.setStatus(PipelineStatusEnum.FAIL.getCode());
pipelineInstanceMapper.updateById(pipelineInstanceTb);
}
}
pipelineInstanceTb.setStatus(pipelineStatusEnum.getCode());
pipelineInstanceMapper.updateById(pipelineInstanceTb);
pipelineInstanceDAO.unLock(pipelineInstanceTb);
} catch (Exception e) {
log.error("流水线运行失败", e);
pipelineInstanceTb.setStatus(PipelineStatusEnum.FAIL.getCode());
pipelineInstanceMapper.updateById(pipelineInstanceTb);
pipelineInstanceDAO.unLock(pipelineInstanceTb);
// pending、running、success、fail
List<String> canRestartStatus = new ArrayList<>() {{
add(PipelineStatusEnum.SUCCESS.getCode());
add(PipelineStatusEnum.FAIL.getCode());
}};
if (!canRestartStatus.contains(currentInstance.getStatus())) {
throw new BadRequestException("流水线运行中,无法重复执行");
}
log.info("解锁流水线, {}", JSON.toJSONString(currentInstance));
pipelineInstanceDAO.unLock(currentInstance);
pipelineJobManage.startJobByNodeCode(currentInstance, retryNodeCode);
}
public void restart(PipelineInstanceTb pipelineInstanceTb) {
boolean exist = pipelineInstanceDAO.lock(pipelineInstanceTb);
if (exist) {
if (!PipelineStatusEnum.FAIL.getCode().equals(pipelineInstanceTb.getCurrentNodeStatus())) {
throw new BadRequestException("流水线执行中,无法重启流水线");
@Override
public List<PipelineNodeDataVo> queryLastPipelineDetail(String instanceIdStr) {
String realInstanceIdStr = instanceIdStr.replace(PipelineConst.INSTANCE_ID, "");
Long instanceId = Long.valueOf(realInstanceIdStr);
List<PipelineNodeDataVo> records = new ArrayList<>();
PipelineInstanceTb pipelineInstanceTb = pipelineInstanceMapper.selectById(instanceId);
if (pipelineInstanceTb == null) {
throw new BadRequestException("流水线实例不存在");
}
log.info("解锁流水线, {}", JSON.toJSONString(pipelineInstanceTb));
pipelineInstanceDAO.unLock(pipelineInstanceTb);
String status = pipelineInstanceTb.getStatus();
List<PipelineInstanceNodeTb> instanceNodeTbs = pipelineInstanceNodeService.queryPipelineInstanceNodeListByInstanceId(instanceId);
for (PipelineInstanceNodeTb instanceNodeTb : instanceNodeTbs) {
PipelineNodeDataVo dataVo = BeanUtil.copyProperties(instanceNodeTb, PipelineNodeDataVo.class);
dataVo.setStatus(status);
// 取流水线节点最后一步的信息
PipelineInstanceNodeDetailTb pipelineInstanceNodeDetail = pipelineInstanceNodeDetailService.getLastPipelineInstanceNodeDetailByArgs(instanceNodeTb, instanceNodeTb.getId(), instanceNodeTb.getCode());
dataVo.setCurrentNodeMsg(pipelineInstanceNodeDetail.getStepMsg());
dataVo.setCurrentNodeStatus(pipelineInstanceNodeDetail.getStepStatus());
records.add(dataVo);
}
return records;
}
}
@@ -20,5 +20,6 @@ import java.util.List;
@Mapper
public interface SystemOssStorageMapper extends BaseMapper<SystemOssStorageTb> {
IPage<SystemOssStorageTb> selectOssStorageByArgs(SystemQueryStorageArgs criteria, Page<SystemOssStorageTb> page);
List<SystemOssStorageTb> selectOssStorageByArgs(SystemQueryStorageArgs criteria);
}
@@ -18,7 +18,6 @@ import cn.odboy.util.ClassUtil;
import cn.odboy.util.FileUtil;
import cn.odboy.util.StringUtil;
import lombok.RequiredArgsConstructor;
import org.springframework.context.annotation.Lazy;
import org.springframework.stereotype.Service;
import org.springframework.transaction.annotation.Transactional;
@@ -29,12 +28,11 @@ import java.util.*;
import java.util.stream.Collectors;
@Service
@RequiredArgsConstructor(onConstructor_ = @Lazy)
@RequiredArgsConstructor
public class SystemDeptService {
private final SystemDeptMapper systemDeptMapper;
private final SystemUserMapper systemUserMapper;
private final SystemRoleMapper systemRoleMapper;
private final SystemDeptService systemDeptService;
/**
* 创建
@@ -45,7 +43,7 @@ public class SystemDeptService {
@Transactional(rollbackFor = Exception.class)
public void saveDept(SystemCreateDeptArgs resources) {
systemDeptMapper.insert(BeanUtil.copyProperties(resources, SystemDeptTb.class));
systemDeptService.updateDeptSubCnt(resources.getPid());
this.updateDeptSubCnt(resources.getPid());
}
/**
@@ -65,8 +63,8 @@ public class SystemDeptService {
SystemDeptTb dept = systemDeptMapper.selectById(resources.getId());
resources.setId(dept.getId());
systemDeptMapper.insertOrUpdate(resources);
systemDeptService.updateDeptSubCnt(oldPid);
systemDeptService.updateDeptSubCnt(newPid);
this.updateDeptSubCnt(oldPid);
this.updateDeptSubCnt(newPid);
}
/**
@@ -79,7 +77,7 @@ public class SystemDeptService {
public void removeDeptByIds(Set<SystemDeptTb> deptSet) {
for (SystemDeptTb dept : deptSet) {
systemDeptMapper.deleteById(dept.getId());
systemDeptService.updateDeptSubCnt(dept.getPid());
this.updateDeptSubCnt(dept.getPid());
}
}
@@ -16,7 +16,6 @@ import cn.odboy.util.ClassUtil;
import cn.odboy.util.FileUtil;
import cn.odboy.util.StringUtil;
import lombok.RequiredArgsConstructor;
import org.springframework.context.annotation.Lazy;
import org.springframework.stereotype.Service;
import org.springframework.transaction.annotation.Transactional;
@@ -28,12 +27,11 @@ import java.util.stream.Collectors;
@Service
@RequiredArgsConstructor(onConstructor_ = @Lazy)
@RequiredArgsConstructor
public class SystemMenuService {
private final SystemMenuMapper systemMenuMapper;
private final SystemRoleMenuMapper systemRoleMenuMapper;
private final SystemRoleService systemRoleService;
private final SystemMenuService systemMenuService;
private static final String YES_STR = "是";
private static final String NO_STR = "否";
@@ -65,7 +63,7 @@ public class SystemMenuService {
// 计算子节点数目
resources.setSubCount(0);
// 更新父节点菜单数目
systemMenuService.updateMenuSubCnt(resources.getPid());
this.updateMenuSubCnt(resources.getPid());
}
/**
@@ -122,8 +120,8 @@ public class SystemMenuService {
menu.setUpdateTime(null);
systemMenuMapper.insertOrUpdate(menu);
// 计算父级菜单节点数目
systemMenuService.updateMenuSubCnt(oldPid);
systemMenuService.updateMenuSubCnt(newPid);
this.updateMenuSubCnt(oldPid);
this.updateMenuSubCnt(newPid);
}
/**
@@ -137,7 +135,7 @@ public class SystemMenuService {
for (SystemMenuTb menu : menuSet) {
systemRoleMenuMapper.deleteRoleMenuByMenuId(menu.getId());
systemMenuMapper.deleteById(menu.getId());
systemMenuService.updateMenuSubCnt(menu.getPid());
this.updateMenuSubCnt(menu.getPid());
}
}
+1 -1
View File
@@ -8,7 +8,7 @@ services:
volumes: # 数据卷挂载路径设置,将本机目录映射到容器目录
- "./mysql/my.cnf:/etc/mysql/my.cnf"
- "./mysql/data:/var/lib/mysql"
# - "./mysql/conf.d:/etc/mysql/conf.d"
# - "./mysql/conf.d:/etc/mysql/conf.d"
- "./mysql/mysql-files:/var/lib/mysql-files"
environment: # 设置环境变量,相当于docker run命令中的-e
TZ: Asia/Shanghai
+1 -1
View File
@@ -355,7 +355,7 @@
<dependency>
<groupId>com.gitee.anwena</groupId>
<artifactId>mongo-plus-boot-starter</artifactId>
<!-- <version>2.0.9.2</version>-->
<!-- <version>2.0.9.2</version>-->
<version>2.1.6.1</version>
</dependency>
<dependency>