From e23711e340070f89895d73c35489c3a6aedeff98 Mon Sep 17 00:00:00 2001 From: Odboy Date: Wed, 23 Jul 2025 23:09:59 +0800 Subject: [PATCH] =?UTF-8?q?feat(devops):=20=E5=AE=9E=E7=8E=B0=E6=B5=81?= =?UTF-8?q?=E6=B0=B4=E7=BA=BF=E7=8A=B6=E6=80=81=E5=AE=9E=E6=97=B6=E6=8E=A8?= =?UTF-8?q?=E9=80=81=E5=8A=9F=E8=83=BD?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - 新增 WebSocket 实时推送流水线状态功能 - 优化流水线数据获取方式,改为实时推送 - 重构 WebSocket 相关代码,提高稳定性和可维护性 - 新增前端 WebSocket 客户端,用于接收实时推送数据 - 修改前端流水线页面,支持实时数据更新 --- .../src/api/devops/pipelineInstance.js | 10 +++ cutejava-front/src/utils/CsWsClient.js | 25 ++++--- .../componentsDemo/CutePipelineNodeDemo.vue | 72 +++++++++++++++++-- .../websocket/constant/CsWsBizTypeEnum.java | 21 ------ .../websocket/context/CsWsClientManager.java | 39 ++++++++++ .../websocket/context/CsWsMessage.java | 18 ++--- .../websocket/context/CsWsServer.java | 39 +++++++--- .../websocket/context/WsClientManager.java | 33 --------- cutejava/cutejava-module-devops/pom.xml | 2 +- .../DevopsPipelineInstanceController.java | 8 +++ .../job/biz/PipelineNodeDeployJavaBiz.java | 2 +- .../devops/job/biz/PipelineNodeInitBiz.java | 8 +-- .../job/biz/PipelineNodeMergeBranchBiz.java | 6 +- .../service/PipelineInstanceService.java | 7 +- .../impl/PipelineInstanceServiceImpl.java | 38 ++++++++++ .../src/main/resources/application-dev.yml | 2 +- 16 files changed, 224 insertions(+), 106 deletions(-) delete mode 100644 cutejava/cutejava-framework/src/main/java/cn/odboy/framework/websocket/constant/CsWsBizTypeEnum.java create mode 100644 cutejava/cutejava-framework/src/main/java/cn/odboy/framework/websocket/context/CsWsClientManager.java delete mode 100644 cutejava/cutejava-framework/src/main/java/cn/odboy/framework/websocket/context/WsClientManager.java diff --git a/cutejava-front/src/api/devops/pipelineInstance.js b/cutejava-front/src/api/devops/pipelineInstance.js index 396c70f6..c0b90cac 100644 --- a/cutejava-front/src/api/devops/pipelineInstance.js +++ b/cutejava-front/src/api/devops/pipelineInstance.js @@ -23,3 +23,13 @@ export function queryLastPipelineDetail(instanceId) { } }) } + +export function queryLastPipelineDetailWs(instanceId) { + return request({ + url: 'api/devops/pipelineInstance/lastWs', + method: 'post', + data: { + instanceId + } + }) +} diff --git a/cutejava-front/src/utils/CsWsClient.js b/cutejava-front/src/utils/CsWsClient.js index 0f36c2cd..444a136b 100644 --- a/cutejava-front/src/utils/CsWsClient.js +++ b/cutejava-front/src/utils/CsWsClient.js @@ -1,25 +1,22 @@ import store from '@/store' -import Vue from 'vue' /** * WebSocket客户端 */ -export default function() { - const user = store.getters && store.getters.user - const webSocketApi = store.getters && store.getters.websocketApi - this.context = Vue.prototype - this.wsUri = webSocketApi.replace('{sid}', user.username) +function CsWsClient(sid) { + this.wsUri = store.getters && store.getters.websocketApi.replace('{sid}', sid) + this.client = null /** * 初始化WebSocket连接 * @param handleError -> function (e) { e是异常本身 } * @param handleMessage -> function (e) { e.data是message } */ 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 } /** @@ -27,16 +24,18 @@ export default function() { * @param data 客户端数据 */ this.sendData = function(data) { - this.wsClient.send(JSON.stringify(data)) + this.client.send(JSON.stringify(data)) } /** * 用于关闭ws连接 */ this.close = function() { - console.log('关闭连接') - if (this.wsClient.readyState === 1) { + console.info('关闭连接') + if (this.client.readyState === 1) { this.sendData([{}]) - this.wsClient.close() + this.client.close() } } } + +export default CsWsClient diff --git a/cutejava-front/src/views/componentsDemo/CutePipelineNodeDemo.vue b/cutejava-front/src/views/componentsDemo/CutePipelineNodeDemo.vue index ef9f2817..ddba5443 100644 --- a/cutejava-front/src/views/componentsDemo/CutePipelineNodeDemo.vue +++ b/cutejava-front/src/views/componentsDemo/CutePipelineNodeDemo.vue @@ -67,9 +67,10 @@ import CutePipelineNode from '@/views/components/dev/CutePipelineNode' 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 CsMessage from '@/utils/elementui/CsMessage' +import CsWsClient from '@/utils/CsWsClient' export default { name: 'CutePipelineNodeDemo', @@ -458,16 +459,18 @@ export default { dynamicStartButtonLoading: false, dynamicTemplate: [], dynamicHook: null, - dynamicInstance: {} + dynamicInstance: {}, + dynamicWsClient: null } }, mounted() { this.refreshData() this.initPipelineTemplate(4) - const pipelineInstanceId = sessionStorage.getItem('pipelineInstanceId') - if (pipelineInstanceId) { - this.fetchLastDetail(pipelineInstanceId) - } + // const pipelineInstanceId = sessionStorage.getItem('pipelineInstanceId') + // if (pipelineInstanceId) { + // this.fetchLastDetail(pipelineInstanceId) + // } + // this.connectWebSocketServer(pipelineInstanceId) }, methods: { refreshData() { @@ -486,6 +489,11 @@ export default { } this.dataList3 = [...this.dataList3] }, + /** + * 初始化流水线模板 + * @param templateId + * @returns {Promise} + */ async initPipelineTemplate(templateId) { const pipelineTemplate = await getPipelineTemplate(templateId) 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} + */ async fetchLastDetail(instanceId) { const that = this const intervalTime = 2000 @@ -517,6 +570,10 @@ export default { that.dynamicStartButtonLoading = false }, 2000) }, + /** + * 启动流水线 + * @returns {Promise} + */ async startPipelineTest() { const that = this try { @@ -535,8 +592,9 @@ export default { // } const result = await startPipeline() CsMessage.Success('流水线启动成功') - await that.fetchLastDetail(result.instanceId) + // await that.fetchLastDetail(result.instanceId) sessionStorage.setItem('pipelineInstanceId', result.instanceId) + this.connectWebSocketServer(result.instanceId) } catch (e) { that.dynamicStartButtonLoading = false } diff --git a/cutejava/cutejava-framework/src/main/java/cn/odboy/framework/websocket/constant/CsWsBizTypeEnum.java b/cutejava/cutejava-framework/src/main/java/cn/odboy/framework/websocket/constant/CsWsBizTypeEnum.java deleted file mode 100644 index 5f9cfe26..00000000 --- a/cutejava/cutejava-framework/src/main/java/cn/odboy/framework/websocket/constant/CsWsBizTypeEnum.java +++ /dev/null @@ -1,21 +0,0 @@ -package cn.odboy.framework.websocket.constant; - - -public enum CsWsBizTypeEnum { - /** - * 连接 - */ - CONNECT, - /** - * 关闭 - */ - CLOSE, - /** - * 信息 - */ - INFO, - /** - * 错误 - */ - ERROR -} diff --git a/cutejava/cutejava-framework/src/main/java/cn/odboy/framework/websocket/context/CsWsClientManager.java b/cutejava/cutejava-framework/src/main/java/cn/odboy/framework/websocket/context/CsWsClientManager.java new file mode 100644 index 00000000..f516f686 --- /dev/null +++ b/cutejava/cutejava-framework/src/main/java/cn/odboy/framework/websocket/context/CsWsClientManager.java @@ -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 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 getAllClient() { + return client.values(); + } + + public static CsWsServer getClientBySid(String sid) { + return client.get(sid); + } +} diff --git a/cutejava/cutejava-framework/src/main/java/cn/odboy/framework/websocket/context/CsWsMessage.java b/cutejava/cutejava-framework/src/main/java/cn/odboy/framework/websocket/context/CsWsMessage.java index 81b3d604..68e11693 100644 --- a/cutejava/cutejava-framework/src/main/java/cn/odboy/framework/websocket/context/CsWsMessage.java +++ b/cutejava/cutejava-framework/src/main/java/cn/odboy/framework/websocket/context/CsWsMessage.java @@ -1,15 +1,17 @@ package cn.odboy.framework.websocket.context; -import cn.odboy.framework.websocket.constant.CsWsBizTypeEnum; -import lombok.Data; +import cn.odboy.base.CsObject; +import lombok.Getter; +import lombok.Setter; -@Data -public class CsWsMessage { - private String message; - private CsWsBizTypeEnum bizType; +@Getter +@Setter +public class CsWsMessage extends CsObject { + private String bizType; + private String data; - public CsWsMessage(String message, CsWsBizTypeEnum bizType) { - this.message = message; + public CsWsMessage(String bizType, String data) { this.bizType = bizType; + this.data = data; } } diff --git a/cutejava/cutejava-framework/src/main/java/cn/odboy/framework/websocket/context/CsWsServer.java b/cutejava/cutejava-framework/src/main/java/cn/odboy/framework/websocket/context/CsWsServer.java index 1f7e4af6..f0ab3dc4 100644 --- a/cutejava/cutejava-framework/src/main/java/cn/odboy/framework/websocket/context/CsWsServer.java +++ b/cutejava/cutejava-framework/src/main/java/cn/odboy/framework/websocket/context/CsWsServer.java @@ -1,8 +1,8 @@ package cn.odboy.framework.websocket.context; import com.alibaba.fastjson2.JSON; +import lombok.Getter; import lombok.extern.slf4j.Slf4j; -import org.springframework.beans.factory.annotation.Autowired; import org.springframework.stereotype.Component; import javax.websocket.*; @@ -15,9 +15,8 @@ import java.util.Objects; @Slf4j @Component @ServerEndpoint("/websocket/{sid}") +@Getter public class CsWsServer { - @Autowired - private WsClientManager wsClientManager; /** * 与某个客户端的连接会话, 需要通过它来给客户端发送数据 */ @@ -34,7 +33,7 @@ public class CsWsServer { public void onOpen(Session session, @PathParam("sid") String sid) { this.session = session; this.sid = sid; - wsClientManager.addClient(sid, this); + CsWsClientManager.addClient(sid, this); } /** @@ -42,7 +41,11 @@ public class CsWsServer { */ @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 public void onMessage(String message, Session session) { - log.info("收到来 sid={} 的信息: {}", sid, message); - sendToAll(message); + CsWsMessage wsMessage = JSON.parseObject(message, CsWsMessage.class); + 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 / */ private void sendToAll(String message) { - for (CsWsServer item : wsClientManager.getAllClient()) { + for (CsWsServer item : CsWsClientManager.getAllClient()) { try { item.innerSendMessage(message); } catch (IOException e) { log.error("发送消息给 sid={} 失败", item.sid, e); + CsWsClientManager.removeClient(item.sid); } } } @@ -74,6 +80,15 @@ public class CsWsServer { @OnError public void onError(Session session, Throwable 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) { sendToAll(body); } else { - wsClientManager.getClientBySid(sid).innerSendMessage(body); + CsWsServer wsClient = CsWsClientManager.getClientBySid(sid); + if (wsClient != null) { + wsClient.innerSendMessage(body); + } } } catch (Exception ignored) { } @@ -108,8 +126,7 @@ public class CsWsServer { return false; } CsWsServer that = (CsWsServer) o; - return Objects.equals(session, that.session) && - Objects.equals(sid, that.sid); + return Objects.equals(session, that.session) && Objects.equals(sid, that.sid); } @Override diff --git a/cutejava/cutejava-framework/src/main/java/cn/odboy/framework/websocket/context/WsClientManager.java b/cutejava/cutejava-framework/src/main/java/cn/odboy/framework/websocket/context/WsClientManager.java deleted file mode 100644 index f0a32dbd..00000000 --- a/cutejava/cutejava-framework/src/main/java/cn/odboy/framework/websocket/context/WsClientManager.java +++ /dev/null @@ -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 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 getAllClient() { - return client.values(); - } - - public CsWsServer getClientBySid(String sid) { - return client.get(sid); - } -} diff --git a/cutejava/cutejava-module-devops/pom.xml b/cutejava/cutejava-module-devops/pom.xml index a10f2964..b05fbd8d 100644 --- a/cutejava/cutejava-module-devops/pom.xml +++ b/cutejava/cutejava-module-devops/pom.xml @@ -13,7 +13,7 @@ cn.odboy - cutejava-framework + cutejava-module-system 1.4.1 diff --git a/cutejava/cutejava-module-devops/src/main/java/cn/odboy/devops/controller/DevopsPipelineInstanceController.java b/cutejava/cutejava-module-devops/src/main/java/cn/odboy/devops/controller/DevopsPipelineInstanceController.java index 3e4af329..aa41d89b 100644 --- a/cutejava/cutejava-module-devops/src/main/java/cn/odboy/devops/controller/DevopsPipelineInstanceController.java +++ b/cutejava/cutejava-module-devops/src/main/java/cn/odboy/devops/controller/DevopsPipelineInstanceController.java @@ -75,4 +75,12 @@ public class DevopsPipelineInstanceController { public ResponseEntity queryLastPipelineDetail(@Validated @RequestBody DevOpsQueryLastPipelineDetailArgs args) { 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("开始推送流水线明细"); + } } diff --git a/cutejava/cutejava-module-devops/src/main/java/cn/odboy/devops/job/biz/PipelineNodeDeployJavaBiz.java b/cutejava/cutejava-module-devops/src/main/java/cn/odboy/devops/job/biz/PipelineNodeDeployJavaBiz.java index 0c8a6c76..bde902f1 100644 --- a/cutejava/cutejava-module-devops/src/main/java/cn/odboy/devops/job/biz/PipelineNodeDeployJavaBiz.java +++ b/cutejava/cutejava-module-devops/src/main/java/cn/odboy/devops/job/biz/PipelineNodeDeployJavaBiz.java @@ -20,7 +20,7 @@ public class PipelineNodeDeployJavaBiz { @PipelineNodeStepLog("部署中") public void deployJavaByWithContextName(PipelineInstanceNodeTb pipelineInstanceNode, String contextName, String env, List templateList, PipelineNodeJobExecuteResult lastNodeResult, String templateType) { - ThreadUtil.safeSleep(10000); + ThreadUtil.safeSleep(2000); } @PipelineNodeStepLog("部署完成") diff --git a/cutejava/cutejava-module-devops/src/main/java/cn/odboy/devops/job/biz/PipelineNodeInitBiz.java b/cutejava/cutejava-module-devops/src/main/java/cn/odboy/devops/job/biz/PipelineNodeInitBiz.java index 0d139b5e..8a0c4919 100644 --- a/cutejava/cutejava-module-devops/src/main/java/cn/odboy/devops/job/biz/PipelineNodeInitBiz.java +++ b/cutejava/cutejava-module-devops/src/main/java/cn/odboy/devops/job/biz/PipelineNodeInitBiz.java @@ -15,21 +15,21 @@ import java.util.List; public class PipelineNodeInitBiz { @PipelineNodeStepLog("初始化开始") public void initStart(PipelineInstanceNodeTb pipelineInstanceNode, String contextName, String env, List templateList, PipelineNodeJobExecuteResult lastNodeResult) { - ThreadUtil.safeSleep(5000); + ThreadUtil.safeSleep(2000); } @PipelineNodeStepLog("Master分支合并到Release分支") public void mergeMasterToRelease(PipelineInstanceNodeTb pipelineInstanceNode, String contextName, String env, List templateList, PipelineNodeJobExecuteResult lastNodeResult) { - ThreadUtil.safeSleep(5000); + ThreadUtil.safeSleep(2000); } @PipelineNodeStepLog("新建Release分支") public void createReleaseBranch(PipelineInstanceNodeTb pipelineInstanceNode, String contextName, String env, List templateList, PipelineNodeJobExecuteResult lastNodeResult) { - ThreadUtil.safeSleep(5000); + ThreadUtil.safeSleep(2000); } @PipelineNodeStepLog("初始化完成") public void initFinish(PipelineInstanceNodeTb pipelineInstanceNode, String contextName, String env, List templateList, PipelineNodeJobExecuteResult lastNodeResult) { - ThreadUtil.safeSleep(5000); + ThreadUtil.safeSleep(2000); } } diff --git a/cutejava/cutejava-module-devops/src/main/java/cn/odboy/devops/job/biz/PipelineNodeMergeBranchBiz.java b/cutejava/cutejava-module-devops/src/main/java/cn/odboy/devops/job/biz/PipelineNodeMergeBranchBiz.java index a59f83f7..78d9118b 100644 --- a/cutejava/cutejava-module-devops/src/main/java/cn/odboy/devops/job/biz/PipelineNodeMergeBranchBiz.java +++ b/cutejava/cutejava-module-devops/src/main/java/cn/odboy/devops/job/biz/PipelineNodeMergeBranchBiz.java @@ -15,16 +15,16 @@ import java.util.List; public class PipelineNodeMergeBranchBiz { @PipelineNodeStepLog("分支合并开始") public void mergeBranchStart(PipelineInstanceNodeTb pipelineInstanceNode, String contextName, String env, List templateList, PipelineNodeJobExecuteResult lastNodeResult) { - ThreadUtil.safeSleep(5000); + ThreadUtil.safeSleep(2000); } @PipelineNodeStepLog("集成区分支合并到release分支") public void integrationAreaBranchMergeRelease(PipelineInstanceNodeTb pipelineInstanceNode, String contextName, String env, List templateList, PipelineNodeJobExecuteResult lastNodeResult) { - ThreadUtil.safeSleep(5000); + ThreadUtil.safeSleep(2000); } @PipelineNodeStepLog("分支合并完成") public void mergeBranchFinish(PipelineInstanceNodeTb pipelineInstanceNode, String contextName, String env, List templateList, PipelineNodeJobExecuteResult lastNodeResult) { - ThreadUtil.safeSleep(5000); + ThreadUtil.safeSleep(2000); } } diff --git a/cutejava/cutejava-module-devops/src/main/java/cn/odboy/devops/service/PipelineInstanceService.java b/cutejava/cutejava-module-devops/src/main/java/cn/odboy/devops/service/PipelineInstanceService.java index edacd69a..93d6da3c 100644 --- a/cutejava/cutejava-module-devops/src/main/java/cn/odboy/devops/service/PipelineInstanceService.java +++ b/cutejava/cutejava-module-devops/src/main/java/cn/odboy/devops/service/PipelineInstanceService.java @@ -3,9 +3,6 @@ package cn.odboy.devops.service; import cn.odboy.devops.dal.dataobject.PipelineInstanceTb; import cn.odboy.devops.dal.model.StartPipelineResultVo; import cn.odboy.devops.framework.pipeline.model.PipelineInstanceVo; -import cn.odboy.devops.framework.pipeline.model.PipelineNodeDataVo; - -import java.util.List; /** *

@@ -17,6 +14,10 @@ import java.util.List; */ public interface PipelineInstanceService { StartPipelineResultVo startPipeline(PipelineInstanceTb pipelineInstanceTb); + StartPipelineResultVo restartPipeline(PipelineInstanceTb pipelineInstanceTb, String retryNodeCode); + PipelineInstanceVo queryLastPipelineDetail(String instanceId); + + void queryLastPipelineDetailWs(String instanceId); } diff --git a/cutejava/cutejava-module-devops/src/main/java/cn/odboy/devops/service/impl/PipelineInstanceServiceImpl.java b/cutejava/cutejava-module-devops/src/main/java/cn/odboy/devops/service/impl/PipelineInstanceServiceImpl.java index a464fcf1..1e3d9ddb 100644 --- a/cutejava/cutejava-module-devops/src/main/java/cn/odboy/devops/service/impl/PipelineInstanceServiceImpl.java +++ b/cutejava/cutejava-module-devops/src/main/java/cn/odboy/devops/service/impl/PipelineInstanceServiceImpl.java @@ -1,6 +1,7 @@ package cn.odboy.devops.service.impl; 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.PipelineStatusEnum; 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.PipelineInstanceService; 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 lombok.RequiredArgsConstructor; import lombok.extern.slf4j.Slf4j; import org.springframework.stereotype.Service; +import java.io.IOException; import java.util.ArrayList; +import java.util.HashMap; import java.util.List; +import java.util.Map; @Slf4j @Service @@ -33,6 +39,8 @@ public class PipelineInstanceServiceImpl implements PipelineInstanceService { private final PipelineInstanceDAO pipelineInstanceDAO; private final PipelineInstanceNodeService pipelineInstanceNodeService; private final PipelineInstanceNodeDetailService pipelineInstanceNodeDetailService; + private final CsWsServer csWsServer; + private volatile Map runningThreadMap = new HashMap<>(); @Override public StartPipelineResultVo startPipeline(PipelineInstanceTb pipelineInstanceTb) { @@ -99,4 +107,34 @@ public class PipelineInstanceServiceImpl implements PipelineInstanceService { pipelineInstanceVo.setNodes(records); 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); + } } diff --git a/cutejava/cutejava-starter/src/main/resources/application-dev.yml b/cutejava/cutejava-starter/src/main/resources/application-dev.yml index caed05dd..90ea1881 100644 --- a/cutejava/cutejava-starter/src/main/resources/application-dev.yml +++ b/cutejava/cutejava-starter/src/main/resources/application-dev.yml @@ -237,7 +237,7 @@ mongo-plus: password: 123456 authentication-database: cutejava log: true - pretty: true + pretty: false configuration: field: ignoring-null: false