feat(devops): 实现流水线状态实时推送功能

- 新增 WebSocket 实时推送流水线状态功能
- 优化流水线数据获取方式,改为实时推送
- 重构 WebSocket 相关代码,提高稳定性和可维护性
- 新增前端 WebSocket 客户端,用于接收实时推送数据
- 修改前端流水线页面,支持实时数据更新
This commit is contained in:
2025-07-23 23:09:59 +08:00
parent 902173a9e9
commit e23711e340
16 changed files with 224 additions and 106 deletions
@@ -23,3 +23,13 @@ export function queryLastPipelineDetail(instanceId) {
} }
}) })
} }
export function queryLastPipelineDetailWs(instanceId) {
return request({
url: 'api/devops/pipelineInstance/lastWs',
method: 'post',
data: {
instanceId
}
})
}
+12 -13
View File
@@ -1,25 +1,22 @@
import store from '@/store' import store from '@/store'
import Vue from 'vue'
/** /**
* WebSocket客户端 * WebSocket客户端
*/ */
export default function() { function CsWsClient(sid) {
const user = store.getters && store.getters.user this.wsUri = store.getters && store.getters.websocketApi.replace('{sid}', sid)
const webSocketApi = store.getters && store.getters.websocketApi this.client = null
this.context = Vue.prototype
this.wsUri = webSocketApi.replace('{sid}', user.username)
/** /**
* 初始化WebSocket连接 * 初始化WebSocket连接
* @param handleError -> function (e) { e是异常本身 } * @param handleError -> function (e) { e是异常本身 }
* @param handleMessage -> function (e) { e.data是message } * @param handleMessage -> function (e) { e.data是message }
*/ */
this.connect = function(handleError, handleMessage) { this.connect = function(handleError, handleMessage) {
this.wsClient = new WebSocket(this.wsUri) this.client = new WebSocket(this.wsUri)
// 连接发生错误 // 连接发生错误
this.wsClient.onerror = handleError this.client.onerror = handleError
// 收到消息 // 收到消息
this.wsClient.onmessage = handleMessage this.client.onmessage = handleMessage
return this return this
} }
/** /**
@@ -27,16 +24,18 @@ export default function() {
* @param data 客户端数据 * @param data 客户端数据
*/ */
this.sendData = function(data) { this.sendData = function(data) {
this.wsClient.send(JSON.stringify(data)) this.client.send(JSON.stringify(data))
} }
/** /**
* 用于关闭ws连接 * 用于关闭ws连接
*/ */
this.close = function() { this.close = function() {
console.log('关闭连接') console.info('关闭连接')
if (this.wsClient.readyState === 1) { if (this.client.readyState === 1) {
this.sendData([{}]) this.sendData([{}])
this.wsClient.close() this.client.close()
} }
} }
} }
export default CsWsClient
@@ -67,9 +67,10 @@
import CutePipelineNode from '@/views/components/dev/CutePipelineNode' import CutePipelineNode from '@/views/components/dev/CutePipelineNode'
import { getPipelineTemplate } from '@/api/devops/pipelineTemplate' import { getPipelineTemplate } from '@/api/devops/pipelineTemplate'
import { queryLastPipelineDetail, restartPipeline, startPipeline } from '@/api/devops/pipelineInstance' import { queryLastPipelineDetail, queryLastPipelineDetailWs, startPipeline } from '@/api/devops/pipelineInstance'
import { CountArraysObjectByPropKey, FormatDateTimeStr } from '@/utils/CsUtil' import { CountArraysObjectByPropKey, FormatDateTimeStr } from '@/utils/CsUtil'
import CsMessage from '@/utils/elementui/CsMessage' import CsMessage from '@/utils/elementui/CsMessage'
import CsWsClient from '@/utils/CsWsClient'
export default { export default {
name: 'CutePipelineNodeDemo', name: 'CutePipelineNodeDemo',
@@ -458,16 +459,18 @@ export default {
dynamicStartButtonLoading: false, dynamicStartButtonLoading: false,
dynamicTemplate: [], dynamicTemplate: [],
dynamicHook: null, dynamicHook: null,
dynamicInstance: {} dynamicInstance: {},
dynamicWsClient: null
} }
}, },
mounted() { mounted() {
this.refreshData() this.refreshData()
this.initPipelineTemplate(4) this.initPipelineTemplate(4)
const pipelineInstanceId = sessionStorage.getItem('pipelineInstanceId') // const pipelineInstanceId = sessionStorage.getItem('pipelineInstanceId')
if (pipelineInstanceId) { // if (pipelineInstanceId) {
this.fetchLastDetail(pipelineInstanceId) // this.fetchLastDetail(pipelineInstanceId)
} // }
// this.connectWebSocketServer(pipelineInstanceId)
}, },
methods: { methods: {
refreshData() { refreshData() {
@@ -486,6 +489,11 @@ export default {
} }
this.dataList3 = [...this.dataList3] this.dataList3 = [...this.dataList3]
}, },
/**
* 初始化流水线模板
* @param templateId
* @returns {Promise<void>}
*/
async initPipelineTemplate(templateId) { async initPipelineTemplate(templateId) {
const pipelineTemplate = await getPipelineTemplate(templateId) const pipelineTemplate = await getPipelineTemplate(templateId)
if (pipelineTemplate && pipelineTemplate.template) { if (pipelineTemplate && pipelineTemplate.template) {
@@ -496,6 +504,51 @@ export default {
} }
} }
}, },
/**
* 连接WebSocket服务
*/
connectWebSocketServer(instanceId) {
const that = this
const sid = 'instanceId_' + instanceId
that.dynamicWsClient = new CsWsClient(sid)
that.dynamicWsClient.connect(that.handleWebSocketError, that.handleWebSocketMessage)
setTimeout(() => {
// 请求推送流水线数据
queryLastPipelineDetailWs(sid)
that.dynamicStartButtonLoading = false
that.dynamicStartupStatus = that.dynamicStartupStatusMap.restart.code
}, 1000)
},
handleWebSocketError(e) {
console.error('连接WebSocket服务失败', e)
},
handleWebSocketMessage(e) {
const that = this
const data = e.data
// console.error('来自服务器的数据:data', data)
const callbackData = JSON.parse(data)
// console.error('callbackData', callbackData)
try {
that.dynamicInstance = JSON.parse(callbackData.data)
// 判断流水线是否结束
const successCount = CountArraysObjectByPropKey(that.dynamicInstance.nodes, 'status', 'success')
if (that.dynamicTemplate && that.dynamicTemplate.length === successCount) {
if (that.dynamicWsClient) {
that.dynamicWsClient.close()
}
that.dynamicStartupStatus = that.dynamicStartupStatusMap.start.code
} else {
that.dynamicStartupStatus = that.dynamicStartupStatusMap.restart.code
}
} catch (e) {
// ignore
}
},
/**
* 通过Http的方式拉取流水线最新的状态
* @param instanceId
* @returns {Promise<void>}
*/
async fetchLastDetail(instanceId) { async fetchLastDetail(instanceId) {
const that = this const that = this
const intervalTime = 2000 const intervalTime = 2000
@@ -517,6 +570,10 @@ export default {
that.dynamicStartButtonLoading = false that.dynamicStartButtonLoading = false
}, 2000) }, 2000)
}, },
/**
* 启动流水线
* @returns {Promise<void>}
*/
async startPipelineTest() { async startPipelineTest() {
const that = this const that = this
try { try {
@@ -535,8 +592,9 @@ export default {
// } // }
const result = await startPipeline() const result = await startPipeline()
CsMessage.Success('流水线启动成功') CsMessage.Success('流水线启动成功')
await that.fetchLastDetail(result.instanceId) // await that.fetchLastDetail(result.instanceId)
sessionStorage.setItem('pipelineInstanceId', result.instanceId) sessionStorage.setItem('pipelineInstanceId', result.instanceId)
this.connectWebSocketServer(result.instanceId)
} catch (e) { } catch (e) {
that.dynamicStartButtonLoading = false that.dynamicStartButtonLoading = false
} }
@@ -1,21 +0,0 @@
package cn.odboy.framework.websocket.constant;
public enum CsWsBizTypeEnum {
/**
* 连接
*/
CONNECT,
/**
* 关闭
*/
CLOSE,
/**
* 信息
*/
INFO,
/**
* 错误
*/
ERROR
}
@@ -0,0 +1,39 @@
package cn.odboy.framework.websocket.context;
import javax.websocket.CloseReason;
import java.util.Collection;
import java.util.Map;
import java.util.concurrent.ConcurrentHashMap;
public class CsWsClientManager {
/**
* concurrent包的线程安全Map, 用来存放每个客户端对应的MyWebSocket对象。(分布式必出问题)
*/
private static final Map<String, CsWsServer> client = new ConcurrentHashMap<>();
public static void addClient(String sid, CsWsServer csWsServer) {
// 如果存在就先删除一个, 防止重复推送消息
removeClient(sid);
client.put(sid, csWsServer);
}
public static void removeClient(String sid) {
try {
CsWsServer wsServer = client.get(sid);
if (wsServer != null) {
wsServer.getSession().close();
}
} catch (Exception e) {
// ignore
}
client.remove(sid);
}
public static Collection<CsWsServer> getAllClient() {
return client.values();
}
public static CsWsServer getClientBySid(String sid) {
return client.get(sid);
}
}
@@ -1,15 +1,17 @@
package cn.odboy.framework.websocket.context; package cn.odboy.framework.websocket.context;
import cn.odboy.framework.websocket.constant.CsWsBizTypeEnum; import cn.odboy.base.CsObject;
import lombok.Data; import lombok.Getter;
import lombok.Setter;
@Data @Getter
public class CsWsMessage { @Setter
private String message; public class CsWsMessage extends CsObject {
private CsWsBizTypeEnum bizType; private String bizType;
private String data;
public CsWsMessage(String message, CsWsBizTypeEnum bizType) { public CsWsMessage(String bizType, String data) {
this.message = message;
this.bizType = bizType; this.bizType = bizType;
this.data = data;
} }
} }
@@ -1,8 +1,8 @@
package cn.odboy.framework.websocket.context; package cn.odboy.framework.websocket.context;
import com.alibaba.fastjson2.JSON; import com.alibaba.fastjson2.JSON;
import lombok.Getter;
import lombok.extern.slf4j.Slf4j; import lombok.extern.slf4j.Slf4j;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.stereotype.Component; import org.springframework.stereotype.Component;
import javax.websocket.*; import javax.websocket.*;
@@ -15,9 +15,8 @@ import java.util.Objects;
@Slf4j @Slf4j
@Component @Component
@ServerEndpoint("/websocket/{sid}") @ServerEndpoint("/websocket/{sid}")
@Getter
public class CsWsServer { public class CsWsServer {
@Autowired
private WsClientManager wsClientManager;
/** /**
* 与某个客户端的连接会话, 需要通过它来给客户端发送数据 * 与某个客户端的连接会话, 需要通过它来给客户端发送数据
*/ */
@@ -34,7 +33,7 @@ public class CsWsServer {
public void onOpen(Session session, @PathParam("sid") String sid) { public void onOpen(Session session, @PathParam("sid") String sid) {
this.session = session; this.session = session;
this.sid = sid; this.sid = sid;
wsClientManager.addClient(sid, this); CsWsClientManager.addClient(sid, this);
} }
/** /**
@@ -42,7 +41,11 @@ public class CsWsServer {
*/ */
@OnClose @OnClose
public void onClose() { public void onClose() {
wsClientManager.removeClient(this.sid); try {
CsWsClientManager.removeClient(this.sid);
} catch (Exception e) {
// ignore
}
} }
/** /**
@@ -52,8 +55,10 @@ public class CsWsServer {
*/ */
@OnMessage @OnMessage
public void onMessage(String message, Session session) { public void onMessage(String message, Session session) {
log.info("收到来 sid={} 的信息: {}", sid, message); CsWsMessage wsMessage = JSON.parseObject(message, CsWsMessage.class);
sendToAll(message); String bizType = wsMessage.getBizType();
Object data = wsMessage.getData();
log.info("收到来 sid={} 的信息: message={}, bizType={}, data={}", sid, message, bizType, JSON.toJSONString(data));
} }
/** /**
@@ -62,11 +67,12 @@ public class CsWsServer {
* @param message / * @param message /
*/ */
private void sendToAll(String message) { private void sendToAll(String message) {
for (CsWsServer item : wsClientManager.getAllClient()) { for (CsWsServer item : CsWsClientManager.getAllClient()) {
try { try {
item.innerSendMessage(message); item.innerSendMessage(message);
} catch (IOException e) { } catch (IOException e) {
log.error("发送消息给 sid={} 失败", item.sid, e); log.error("发送消息给 sid={} 失败", item.sid, e);
CsWsClientManager.removeClient(item.sid);
} }
} }
} }
@@ -74,6 +80,15 @@ public class CsWsServer {
@OnError @OnError
public void onError(Session session, Throwable error) { public void onError(Session session, Throwable error) {
log.error("WebSocket sid={} 发生错误", this.sid, error); log.error("WebSocket sid={} 发生错误", this.sid, error);
CsWsServer client = CsWsClientManager.getClientBySid(this.sid);
if (client != null) {
try {
client.getSession().close();
} catch (IOException e) {
// ignore
}
CsWsClientManager.removeClient(this.sid);
}
} }
/** /**
@@ -93,7 +108,10 @@ public class CsWsServer {
if (sid == null) { if (sid == null) {
sendToAll(body); sendToAll(body);
} else { } else {
wsClientManager.getClientBySid(sid).innerSendMessage(body); CsWsServer wsClient = CsWsClientManager.getClientBySid(sid);
if (wsClient != null) {
wsClient.innerSendMessage(body);
}
} }
} catch (Exception ignored) { } catch (Exception ignored) {
} }
@@ -108,8 +126,7 @@ public class CsWsServer {
return false; return false;
} }
CsWsServer that = (CsWsServer) o; CsWsServer that = (CsWsServer) o;
return Objects.equals(session, that.session) && return Objects.equals(session, that.session) && Objects.equals(sid, that.sid);
Objects.equals(sid, that.sid);
} }
@Override @Override
@@ -1,33 +0,0 @@
package cn.odboy.framework.websocket.context;
import org.springframework.stereotype.Component;
import java.util.Collection;
import java.util.Map;
import java.util.concurrent.ConcurrentHashMap;
@Component
public class WsClientManager {
/**
* concurrent包的线程安全Map, 用来存放每个客户端对应的MyWebSocket对象。(分布式必出问题)
*/
private final Map<String, CsWsServer> client = new ConcurrentHashMap<>();
public void addClient(String sid, CsWsServer csWsServer) {
// 如果存在就先删除一个, 防止重复推送消息
client.remove(sid);
client.put(sid, csWsServer);
}
public void removeClient(String sid) {
client.remove(sid);
}
public Collection<CsWsServer> getAllClient() {
return client.values();
}
public CsWsServer getClientBySid(String sid) {
return client.get(sid);
}
}
+1 -1
View File
@@ -13,7 +13,7 @@
<dependencies> <dependencies>
<dependency> <dependency>
<groupId>cn.odboy</groupId> <groupId>cn.odboy</groupId>
<artifactId>cutejava-framework</artifactId> <artifactId>cutejava-module-system</artifactId>
<version>1.4.1</version> <version>1.4.1</version>
</dependency> </dependency>
</dependencies> </dependencies>
@@ -75,4 +75,12 @@ public class DevopsPipelineInstanceController {
public ResponseEntity<?> queryLastPipelineDetail(@Validated @RequestBody DevOpsQueryLastPipelineDetailArgs args) { public ResponseEntity<?> queryLastPipelineDetail(@Validated @RequestBody DevOpsQueryLastPipelineDetailArgs args) {
return ResponseEntity.ok(pipelineInstanceService.queryLastPipelineDetail(args.getInstanceId())); return ResponseEntity.ok(pipelineInstanceService.queryLastPipelineDetail(args.getInstanceId()));
} }
@ApiOperation("查询流水线明细Ws")
@PostMapping(value = "/lastWs")
@PreAuthorize("@el.check()")
public ResponseEntity<?> queryLastPipelineDetailWs(@Validated @RequestBody DevOpsQueryLastPipelineDetailArgs args) {
pipelineInstanceService.queryLastPipelineDetailWs(args.getInstanceId());
return ResponseEntity.ok("开始推送流水线明细");
}
} }
@@ -20,7 +20,7 @@ public class PipelineNodeDeployJavaBiz {
@PipelineNodeStepLog("部署中") @PipelineNodeStepLog("部署中")
public void deployJavaByWithContextName(PipelineInstanceNodeTb pipelineInstanceNode, String contextName, String env, List<PipelineNodeTemplateVo> templateList, PipelineNodeJobExecuteResult lastNodeResult, String templateType) { public void deployJavaByWithContextName(PipelineInstanceNodeTb pipelineInstanceNode, String contextName, String env, List<PipelineNodeTemplateVo> templateList, PipelineNodeJobExecuteResult lastNodeResult, String templateType) {
ThreadUtil.safeSleep(10000); ThreadUtil.safeSleep(2000);
} }
@PipelineNodeStepLog("部署完成") @PipelineNodeStepLog("部署完成")
@@ -15,21 +15,21 @@ import java.util.List;
public class PipelineNodeInitBiz { public class PipelineNodeInitBiz {
@PipelineNodeStepLog("初始化开始") @PipelineNodeStepLog("初始化开始")
public void initStart(PipelineInstanceNodeTb pipelineInstanceNode, String contextName, String env, List<PipelineNodeTemplateVo> templateList, PipelineNodeJobExecuteResult lastNodeResult) { public void initStart(PipelineInstanceNodeTb pipelineInstanceNode, String contextName, String env, List<PipelineNodeTemplateVo> templateList, PipelineNodeJobExecuteResult lastNodeResult) {
ThreadUtil.safeSleep(5000); ThreadUtil.safeSleep(2000);
} }
@PipelineNodeStepLog("Master分支合并到Release分支") @PipelineNodeStepLog("Master分支合并到Release分支")
public void mergeMasterToRelease(PipelineInstanceNodeTb pipelineInstanceNode, String contextName, String env, List<PipelineNodeTemplateVo> templateList, PipelineNodeJobExecuteResult lastNodeResult) { public void mergeMasterToRelease(PipelineInstanceNodeTb pipelineInstanceNode, String contextName, String env, List<PipelineNodeTemplateVo> templateList, PipelineNodeJobExecuteResult lastNodeResult) {
ThreadUtil.safeSleep(5000); ThreadUtil.safeSleep(2000);
} }
@PipelineNodeStepLog("新建Release分支") @PipelineNodeStepLog("新建Release分支")
public void createReleaseBranch(PipelineInstanceNodeTb pipelineInstanceNode, String contextName, String env, List<PipelineNodeTemplateVo> templateList, PipelineNodeJobExecuteResult lastNodeResult) { public void createReleaseBranch(PipelineInstanceNodeTb pipelineInstanceNode, String contextName, String env, List<PipelineNodeTemplateVo> templateList, PipelineNodeJobExecuteResult lastNodeResult) {
ThreadUtil.safeSleep(5000); ThreadUtil.safeSleep(2000);
} }
@PipelineNodeStepLog("初始化完成") @PipelineNodeStepLog("初始化完成")
public void initFinish(PipelineInstanceNodeTb pipelineInstanceNode, String contextName, String env, List<PipelineNodeTemplateVo> templateList, PipelineNodeJobExecuteResult lastNodeResult) { public void initFinish(PipelineInstanceNodeTb pipelineInstanceNode, String contextName, String env, List<PipelineNodeTemplateVo> templateList, PipelineNodeJobExecuteResult lastNodeResult) {
ThreadUtil.safeSleep(5000); ThreadUtil.safeSleep(2000);
} }
} }
@@ -15,16 +15,16 @@ import java.util.List;
public class PipelineNodeMergeBranchBiz { public class PipelineNodeMergeBranchBiz {
@PipelineNodeStepLog("分支合并开始") @PipelineNodeStepLog("分支合并开始")
public void mergeBranchStart(PipelineInstanceNodeTb pipelineInstanceNode, String contextName, String env, List<PipelineNodeTemplateVo> templateList, PipelineNodeJobExecuteResult lastNodeResult) { public void mergeBranchStart(PipelineInstanceNodeTb pipelineInstanceNode, String contextName, String env, List<PipelineNodeTemplateVo> templateList, PipelineNodeJobExecuteResult lastNodeResult) {
ThreadUtil.safeSleep(5000); ThreadUtil.safeSleep(2000);
} }
@PipelineNodeStepLog("集成区分支合并到release分支") @PipelineNodeStepLog("集成区分支合并到release分支")
public void integrationAreaBranchMergeRelease(PipelineInstanceNodeTb pipelineInstanceNode, String contextName, String env, List<PipelineNodeTemplateVo> templateList, PipelineNodeJobExecuteResult lastNodeResult) { public void integrationAreaBranchMergeRelease(PipelineInstanceNodeTb pipelineInstanceNode, String contextName, String env, List<PipelineNodeTemplateVo> templateList, PipelineNodeJobExecuteResult lastNodeResult) {
ThreadUtil.safeSleep(5000); ThreadUtil.safeSleep(2000);
} }
@PipelineNodeStepLog("分支合并完成") @PipelineNodeStepLog("分支合并完成")
public void mergeBranchFinish(PipelineInstanceNodeTb pipelineInstanceNode, String contextName, String env, List<PipelineNodeTemplateVo> templateList, PipelineNodeJobExecuteResult lastNodeResult) { public void mergeBranchFinish(PipelineInstanceNodeTb pipelineInstanceNode, String contextName, String env, List<PipelineNodeTemplateVo> templateList, PipelineNodeJobExecuteResult lastNodeResult) {
ThreadUtil.safeSleep(5000); ThreadUtil.safeSleep(2000);
} }
} }
@@ -3,9 +3,6 @@ package cn.odboy.devops.service;
import cn.odboy.devops.dal.dataobject.PipelineInstanceTb; import cn.odboy.devops.dal.dataobject.PipelineInstanceTb;
import cn.odboy.devops.dal.model.StartPipelineResultVo; import cn.odboy.devops.dal.model.StartPipelineResultVo;
import cn.odboy.devops.framework.pipeline.model.PipelineInstanceVo; import cn.odboy.devops.framework.pipeline.model.PipelineInstanceVo;
import cn.odboy.devops.framework.pipeline.model.PipelineNodeDataVo;
import java.util.List;
/** /**
* <p> * <p>
@@ -17,6 +14,10 @@ import java.util.List;
*/ */
public interface PipelineInstanceService { public interface PipelineInstanceService {
StartPipelineResultVo startPipeline(PipelineInstanceTb pipelineInstanceTb); StartPipelineResultVo startPipeline(PipelineInstanceTb pipelineInstanceTb);
StartPipelineResultVo restartPipeline(PipelineInstanceTb pipelineInstanceTb, String retryNodeCode); StartPipelineResultVo restartPipeline(PipelineInstanceTb pipelineInstanceTb, String retryNodeCode);
PipelineInstanceVo queryLastPipelineDetail(String instanceId); PipelineInstanceVo queryLastPipelineDetail(String instanceId);
void queryLastPipelineDetailWs(String instanceId);
} }
@@ -1,6 +1,7 @@
package cn.odboy.devops.service.impl; package cn.odboy.devops.service.impl;
import cn.hutool.core.bean.BeanUtil; import cn.hutool.core.bean.BeanUtil;
import cn.hutool.core.thread.ThreadUtil;
import cn.odboy.devops.constant.pipeline.PipelineConst; import cn.odboy.devops.constant.pipeline.PipelineConst;
import cn.odboy.devops.constant.pipeline.PipelineStatusEnum; import cn.odboy.devops.constant.pipeline.PipelineStatusEnum;
import cn.odboy.devops.dal.dataobject.PipelineInstanceNodeDetailTb; import cn.odboy.devops.dal.dataobject.PipelineInstanceNodeDetailTb;
@@ -16,13 +17,18 @@ import cn.odboy.devops.service.PipelineInstanceNodeDetailService;
import cn.odboy.devops.service.PipelineInstanceNodeService; import cn.odboy.devops.service.PipelineInstanceNodeService;
import cn.odboy.devops.service.PipelineInstanceService; import cn.odboy.devops.service.PipelineInstanceService;
import cn.odboy.framework.exception.BadRequestException; import cn.odboy.framework.exception.BadRequestException;
import cn.odboy.framework.websocket.context.CsWsMessage;
import cn.odboy.framework.websocket.context.CsWsServer;
import com.alibaba.fastjson2.JSON; import com.alibaba.fastjson2.JSON;
import lombok.RequiredArgsConstructor; import lombok.RequiredArgsConstructor;
import lombok.extern.slf4j.Slf4j; import lombok.extern.slf4j.Slf4j;
import org.springframework.stereotype.Service; import org.springframework.stereotype.Service;
import java.io.IOException;
import java.util.ArrayList; import java.util.ArrayList;
import java.util.HashMap;
import java.util.List; import java.util.List;
import java.util.Map;
@Slf4j @Slf4j
@Service @Service
@@ -33,6 +39,8 @@ public class PipelineInstanceServiceImpl implements PipelineInstanceService {
private final PipelineInstanceDAO pipelineInstanceDAO; private final PipelineInstanceDAO pipelineInstanceDAO;
private final PipelineInstanceNodeService pipelineInstanceNodeService; private final PipelineInstanceNodeService pipelineInstanceNodeService;
private final PipelineInstanceNodeDetailService pipelineInstanceNodeDetailService; private final PipelineInstanceNodeDetailService pipelineInstanceNodeDetailService;
private final CsWsServer csWsServer;
private volatile Map<String, Thread> runningThreadMap = new HashMap<>();
@Override @Override
public StartPipelineResultVo startPipeline(PipelineInstanceTb pipelineInstanceTb) { public StartPipelineResultVo startPipeline(PipelineInstanceTb pipelineInstanceTb) {
@@ -99,4 +107,34 @@ public class PipelineInstanceServiceImpl implements PipelineInstanceService {
pipelineInstanceVo.setNodes(records); pipelineInstanceVo.setNodes(records);
return pipelineInstanceVo; return pipelineInstanceVo;
} }
@Override
public void queryLastPipelineDetailWs(String instanceId) {
Thread thread = runningThreadMap.get(instanceId);
if (thread != null && !thread.isInterrupted()) {
try {
thread.stop();
} catch (Exception e) {
// ignore
}
}
thread = new Thread(() -> {
boolean loop = true;
final String currentSid = instanceId;
final String realInstanceId = currentSid.replace("instanceId_", "");
while (loop) {
ThreadUtil.safeSleep(1000);
try {
PipelineInstanceVo pipelineInstanceVo = queryLastPipelineDetail(realInstanceId);
CsWsMessage message = new CsWsMessage("FetchPipelineLastDetail", JSON.toJSONString(pipelineInstanceVo));
csWsServer.sendMessage(message, currentSid);
} catch (IOException e) {
log.error("推送流水线最新数据失败", e);
loop = false;
}
}
});
thread.start();
runningThreadMap.put(instanceId, thread);
}
} }
@@ -237,7 +237,7 @@ mongo-plus:
password: 123456 password: 123456
authentication-database: cutejava authentication-database: cutejava
log: true log: true
pretty: true pretty: false
configuration: configuration:
field: field:
ignoring-null: false ignoring-null: false