refactor(websocket): 重构 WebSocket 相关代码
- 修改类名和包名,提高代码一致性 - 优化部分代码结构,提高可维护性 - 移除未使用的流水线相关代码- 更新异常处理逻辑
This commit is contained in:
-16
@@ -1,16 +0,0 @@
|
|||||||
package cn.odboy.framework.exception;
|
|
||||||
|
|
||||||
import org.springframework.util.StringUtils;
|
|
||||||
|
|
||||||
|
|
||||||
public class EntityExistException extends RuntimeException {
|
|
||||||
|
|
||||||
public EntityExistException(Class clazz, String field, String val) {
|
|
||||||
super(EntityExistException.generateMessage(clazz.getSimpleName(), field, val));
|
|
||||||
}
|
|
||||||
|
|
||||||
private static String generateMessage(String entity, String field, String val) {
|
|
||||||
return StringUtils.capitalize(entity)
|
|
||||||
+ " with " + field + " " + val + " existed";
|
|
||||||
}
|
|
||||||
}
|
|
||||||
-15
@@ -1,15 +0,0 @@
|
|||||||
package cn.odboy.framework.exception;
|
|
||||||
|
|
||||||
import org.springframework.util.StringUtils;
|
|
||||||
|
|
||||||
public class EntityNotFoundException extends RuntimeException {
|
|
||||||
|
|
||||||
public EntityNotFoundException(Class clazz, String field, String val) {
|
|
||||||
super(EntityNotFoundException.generateMessage(clazz.getSimpleName(), field, val));
|
|
||||||
}
|
|
||||||
|
|
||||||
private static String generateMessage(String entity, String field, String val) {
|
|
||||||
return StringUtils.capitalize(entity)
|
|
||||||
+ " with " + field + " " + val + " does not exist";
|
|
||||||
}
|
|
||||||
}
|
|
||||||
-22
@@ -2,8 +2,6 @@ package cn.odboy.framework.exception.handler;
|
|||||||
|
|
||||||
import cn.hutool.core.exceptions.ExceptionUtil;
|
import cn.hutool.core.exceptions.ExceptionUtil;
|
||||||
import cn.odboy.framework.exception.BadRequestException;
|
import cn.odboy.framework.exception.BadRequestException;
|
||||||
import cn.odboy.framework.exception.EntityExistException;
|
|
||||||
import cn.odboy.framework.exception.EntityNotFoundException;
|
|
||||||
import lombok.extern.slf4j.Slf4j;
|
import lombok.extern.slf4j.Slf4j;
|
||||||
import org.springframework.http.HttpStatus;
|
import org.springframework.http.HttpStatus;
|
||||||
import org.springframework.http.ResponseEntity;
|
import org.springframework.http.ResponseEntity;
|
||||||
@@ -51,26 +49,6 @@ public class GlobalExceptionHandler {
|
|||||||
return buildResponseEntity(ApiError.error(e.getStatus(), e.getMessage()));
|
return buildResponseEntity(ApiError.error(e.getStatus(), e.getMessage()));
|
||||||
}
|
}
|
||||||
|
|
||||||
/**
|
|
||||||
* 处理 EntityExist
|
|
||||||
*/
|
|
||||||
@ExceptionHandler(value = EntityExistException.class)
|
|
||||||
public ResponseEntity<ApiError> entityExistException(EntityExistException e) {
|
|
||||||
// 打印堆栈信息
|
|
||||||
log.error(ExceptionUtil.stacktraceToString(e));
|
|
||||||
return buildResponseEntity(ApiError.error(e.getMessage()));
|
|
||||||
}
|
|
||||||
|
|
||||||
/**
|
|
||||||
* 处理 EntityNotFound
|
|
||||||
*/
|
|
||||||
@ExceptionHandler(value = EntityNotFoundException.class)
|
|
||||||
public ResponseEntity<ApiError> entityNotFoundException(EntityNotFoundException e) {
|
|
||||||
// 打印堆栈信息
|
|
||||||
log.error(ExceptionUtil.stacktraceToString(e));
|
|
||||||
return buildResponseEntity(ApiError.error(NOT_FOUND.value(), e.getMessage()));
|
|
||||||
}
|
|
||||||
|
|
||||||
/**
|
/**
|
||||||
* 处理所有接口数据验证异常
|
* 处理所有接口数据验证异常
|
||||||
*/
|
*/
|
||||||
|
|||||||
-99
@@ -1,99 +0,0 @@
|
|||||||
//package cn.odboy.system.framework.flow.core;
|
|
||||||
//
|
|
||||||
//import cn.hutool.core.date.DateTime;
|
|
||||||
//import cn.hutool.core.thread.ThreadUtil;
|
|
||||||
//import cn.hutool.core.util.IdUtil;
|
|
||||||
//import cn.odboy.context.SpringBeanHolder;
|
|
||||||
//import cn.odboy.system.framework.flow.core.handler.SerialPipelineNodeHandler;
|
|
||||||
//import cn.odboy.system.framework.flow.model.PipelineNodeVo;
|
|
||||||
//import cn.odboy.exception.BadRequestException;
|
|
||||||
//import lombok.extern.slf4j.Slf4j;
|
|
||||||
//import org.springframework.stereotype.Component;
|
|
||||||
//
|
|
||||||
//import java.util.ArrayList;
|
|
||||||
//import java.util.List;
|
|
||||||
//
|
|
||||||
///**
|
|
||||||
// * 流水线管理工具
|
|
||||||
// *
|
|
||||||
// * @author odboy
|
|
||||||
// */
|
|
||||||
//@Slf4j
|
|
||||||
//@Component
|
|
||||||
//public class PipelineManager {
|
|
||||||
// private static final PipelineTaskPool pipelineTaskPool = new PipelineTaskPool("DefaultPipelineTaskPool");
|
|
||||||
//
|
|
||||||
// public void test() {
|
|
||||||
// String pipelineId = IdUtil.objectId();
|
|
||||||
// List<PipelineNodeVo> pipelineNodeVos = initialize();
|
|
||||||
// execute(pipelineId, pipelineNodeVos);
|
|
||||||
//
|
|
||||||
// ThreadUtil.safeSleep(5000);
|
|
||||||
// pipelineTaskPool.stopTaskForcibly(pipelineId);
|
|
||||||
// ThreadUtil.safeSleep(1000000);
|
|
||||||
// }
|
|
||||||
//
|
|
||||||
// public void execute(String pipelineId, List<PipelineNodeVo> pipelineNodeVos) {
|
|
||||||
// pipelineTaskPool.submitTask(pipelineId, () -> {
|
|
||||||
// for (PipelineNodeVo pipelineNodeVo : pipelineNodeVos) {
|
|
||||||
// executePipeline(pipelineId, pipelineNodeVo);
|
|
||||||
// // 要么成功,要么失败
|
|
||||||
// boolean interrupted = Thread.currentThread().isInterrupted();
|
|
||||||
// while (!interrupted && "running".equals(pipelineNodeVo.getExecuteStatus())) {
|
|
||||||
// log.info("流水线执行中...");
|
|
||||||
// ThreadUtil.safeSleep(2000);
|
|
||||||
// }
|
|
||||||
// if (interrupted) {
|
|
||||||
// log.info("流水线被中断...");
|
|
||||||
// return;
|
|
||||||
// }
|
|
||||||
// }
|
|
||||||
// });
|
|
||||||
// }
|
|
||||||
//
|
|
||||||
// private void executePipeline(String pipelineId, PipelineNodeVo pipelineNodeVo) {
|
|
||||||
// String pipelineServiceName = String.format("pipeline:%s", pipelineNodeVo.getBizCode());
|
|
||||||
// SerialPipelineNodeHandler<Object, Object> pipelineNodeHandler = SpringBeanHolder.getBean(pipelineServiceName);
|
|
||||||
// try {
|
|
||||||
// Object o = pipelineNodeHandler.doProcess(pipelineId, pipelineNodeVo);
|
|
||||||
// log.info("运行结果, {}", o);
|
|
||||||
// } catch (BadRequestException e) {
|
|
||||||
// log.error("运行流水线异常", e);
|
|
||||||
// if (pipelineTaskPool.isTaskRunning(pipelineId)) {
|
|
||||||
// pipelineTaskPool.stopTaskForcibly(pipelineId);
|
|
||||||
// }
|
|
||||||
// } catch (Exception e) {
|
|
||||||
// log.error("其他异常", e);
|
|
||||||
// if (pipelineTaskPool.isTaskRunning(pipelineId)) {
|
|
||||||
// pipelineTaskPool.stopTaskForcibly(pipelineId);
|
|
||||||
// }
|
|
||||||
// }
|
|
||||||
// }
|
|
||||||
//
|
|
||||||
// public List<PipelineNodeVo> initialize() {
|
|
||||||
// List<PipelineNodeVo> pipelineNodeVos = new ArrayList<>();
|
|
||||||
// PipelineNodeVo createReleaseBranchNodeVo = new PipelineNodeVo();
|
|
||||||
// createReleaseBranchNodeVo.setName("初始化");
|
|
||||||
// createReleaseBranchNodeVo.setClick(true);
|
|
||||||
// createReleaseBranchNodeVo.setButtonList(null);
|
|
||||||
// createReleaseBranchNodeVo.setBizCode("create_release_branch");
|
|
||||||
// createReleaseBranchNodeVo.setStartTimeMillis(DateTime.now().getTime());
|
|
||||||
// createReleaseBranchNodeVo.setDurationMillis(0L);
|
|
||||||
// createReleaseBranchNodeVo.setExecuteStatus("running");
|
|
||||||
// createReleaseBranchNodeVo.setResultDesc(null);
|
|
||||||
// pipelineNodeVos.add(createReleaseBranchNodeVo);
|
|
||||||
//
|
|
||||||
// PipelineNodeVo createReleaseBranchNodeVo1 = new PipelineNodeVo();
|
|
||||||
// createReleaseBranchNodeVo1.setName("代码构建");
|
|
||||||
// createReleaseBranchNodeVo1.setClick(true);
|
|
||||||
// createReleaseBranchNodeVo1.setButtonList(null);
|
|
||||||
// createReleaseBranchNodeVo1.setBizCode("build_code");
|
|
||||||
// createReleaseBranchNodeVo1.setStartTimeMillis(DateTime.now().getTime());
|
|
||||||
// createReleaseBranchNodeVo1.setDurationMillis(0L);
|
|
||||||
// createReleaseBranchNodeVo1.setExecuteStatus("running");
|
|
||||||
// createReleaseBranchNodeVo1.setResultDesc(null);
|
|
||||||
// pipelineNodeVos.add(createReleaseBranchNodeVo1);
|
|
||||||
//
|
|
||||||
// return pipelineNodeVos;
|
|
||||||
// }
|
|
||||||
//}
|
|
||||||
-212
@@ -1,212 +0,0 @@
|
|||||||
//package cn.odboy.system.framework.flow.core;
|
|
||||||
//
|
|
||||||
//import cn.hutool.core.date.DateTime;
|
|
||||||
//import cn.hutool.core.thread.ThreadUtil;
|
|
||||||
//import cn.hutool.core.util.IdUtil;
|
|
||||||
//import cn.odboy.context.SpringBeanHolder;
|
|
||||||
//import cn.odboy.system.framework.flow.core.handler.SerialPipelineNodeHandler;
|
|
||||||
//import cn.odboy.system.framework.flow.model.PipelineNodeVo;
|
|
||||||
//import cn.odboy.exception.BadRequestException;
|
|
||||||
//import lombok.extern.slf4j.Slf4j;
|
|
||||||
//import org.springframework.stereotype.Component;
|
|
||||||
//
|
|
||||||
//import java.util.ArrayList;
|
|
||||||
//import java.util.List;
|
|
||||||
//import java.util.Map;
|
|
||||||
//import java.util.Optional;
|
|
||||||
//import java.util.concurrent.ConcurrentHashMap;
|
|
||||||
//import java.util.concurrent.atomic.AtomicBoolean;
|
|
||||||
//
|
|
||||||
///**
|
|
||||||
// * 流水线进程管理工具
|
|
||||||
// *
|
|
||||||
// * @author odboy
|
|
||||||
// */
|
|
||||||
//@Slf4j
|
|
||||||
//@Component
|
|
||||||
//public class PipelineProcessManager {
|
|
||||||
// // 使用ConcurrentHashMap保证线程安全
|
|
||||||
// private final Map<String, Thread> runningTasks = new ConcurrentHashMap<>();
|
|
||||||
// // 使用AtomicBoolean作为停止标志,保证原子性操作
|
|
||||||
// private final Map<String, AtomicBoolean> stopFlags = new ConcurrentHashMap<>();
|
|
||||||
// // 常量定义
|
|
||||||
// private static final String PIPELINE_THREAD_PREFIX = "pipeline-";
|
|
||||||
// private static final String PIPELINE_SERVICE_PREFIX = "pipeline:";
|
|
||||||
// private static final String RUNNING_STATUS = "running";
|
|
||||||
// private static final long DEFAULT_SLEEP_MS = 2000L;
|
|
||||||
//
|
|
||||||
// /**
|
|
||||||
// * 测试方法
|
|
||||||
// */
|
|
||||||
// public void test() {
|
|
||||||
// String pipelineId = IdUtil.objectId();
|
|
||||||
// List<PipelineNodeVo> pipelineNodes = initialize();
|
|
||||||
// execute(pipelineId, pipelineNodes);
|
|
||||||
//
|
|
||||||
// ThreadUtil.safeSleep(5000);
|
|
||||||
// stopForcibly(pipelineId);
|
|
||||||
// ThreadUtil.safeSleep(1000000);
|
|
||||||
// }
|
|
||||||
//
|
|
||||||
// /**
|
|
||||||
// * 执行流水线
|
|
||||||
// *
|
|
||||||
// * @param pipelineId 流水线ID
|
|
||||||
// * @param pipelineNodes 流水线节点列表
|
|
||||||
// */
|
|
||||||
// public void execute(String pipelineId, List<PipelineNodeVo> pipelineNodes) {
|
|
||||||
// if (runningTasks.containsKey(pipelineId)) {
|
|
||||||
// log.warn("Pipeline [{}] is already running", pipelineId);
|
|
||||||
// return;
|
|
||||||
// }
|
|
||||||
//
|
|
||||||
// stopFlags.put(pipelineId, new AtomicBoolean(false));
|
|
||||||
//
|
|
||||||
// Thread pipelineThread = new Thread(() -> processPipeline(pipelineId, pipelineNodes));
|
|
||||||
// pipelineThread.setDaemon(true);
|
|
||||||
// pipelineThread.setName(PIPELINE_THREAD_PREFIX + pipelineId);
|
|
||||||
//
|
|
||||||
// runningTasks.put(pipelineId, pipelineThread);
|
|
||||||
// pipelineThread.start();
|
|
||||||
// }
|
|
||||||
//
|
|
||||||
// /**
|
|
||||||
// * 处理流水线任务
|
|
||||||
// */
|
|
||||||
// private void processPipeline(String pipelineId, List<PipelineNodeVo> pipelineNodes) {
|
|
||||||
// try {
|
|
||||||
// for (PipelineNodeVo node : pipelineNodes) {
|
|
||||||
// if (shouldStop(pipelineId)) {
|
|
||||||
// log.info("Pipeline [{}] was interrupted", pipelineId);
|
|
||||||
// return;
|
|
||||||
// }
|
|
||||||
//
|
|
||||||
// executePipelineNode(pipelineId, node);
|
|
||||||
//
|
|
||||||
// while (!shouldStop(pipelineId) && RUNNING_STATUS.equals(node.getExecuteStatus())) {
|
|
||||||
// log.debug("Pipeline [{}] is running...", pipelineId);
|
|
||||||
// sleepInterruptibly(DEFAULT_SLEEP_MS, pipelineId);
|
|
||||||
// }
|
|
||||||
//
|
|
||||||
// if (shouldStop(pipelineId)) {
|
|
||||||
// log.info("Pipeline [{}] was interrupted", pipelineId);
|
|
||||||
// return;
|
|
||||||
// }
|
|
||||||
// }
|
|
||||||
// } catch (Exception e) {
|
|
||||||
// log.error("流水线异常", e);
|
|
||||||
// } finally {
|
|
||||||
// cleanup(pipelineId);
|
|
||||||
// }
|
|
||||||
// }
|
|
||||||
//
|
|
||||||
// /**
|
|
||||||
// * 检查是否应该停止
|
|
||||||
// */
|
|
||||||
// private boolean shouldStop(String pipelineId) {
|
|
||||||
// return Optional.ofNullable(stopFlags.get(pipelineId))
|
|
||||||
// .map(AtomicBoolean::get)
|
|
||||||
// .orElse(true);
|
|
||||||
// }
|
|
||||||
//
|
|
||||||
// /**
|
|
||||||
// * 可中断的睡眠
|
|
||||||
// */
|
|
||||||
// private void sleepInterruptibly(long millis, String pipelineId) {
|
|
||||||
// try {
|
|
||||||
// Thread.sleep(millis);
|
|
||||||
// } catch (InterruptedException e) {
|
|
||||||
// log.info("Pipeline [{}] sleep interrupted", pipelineId);
|
|
||||||
// Thread.currentThread().interrupt();
|
|
||||||
// }
|
|
||||||
// }
|
|
||||||
//
|
|
||||||
// /**
|
|
||||||
// * 清理资源
|
|
||||||
// */
|
|
||||||
// private void cleanup(String pipelineId) {
|
|
||||||
// runningTasks.remove(pipelineId);
|
|
||||||
// stopFlags.remove(pipelineId);
|
|
||||||
// log.info("Pipeline [{}] resources cleaned up", pipelineId);
|
|
||||||
// }
|
|
||||||
//
|
|
||||||
// /**
|
|
||||||
// * 检查流水线是否在运行
|
|
||||||
// */
|
|
||||||
// public boolean isRunning(String pipelineId) {
|
|
||||||
// return Optional.ofNullable(runningTasks.get(pipelineId))
|
|
||||||
// .map(Thread::isAlive)
|
|
||||||
// .orElse(false);
|
|
||||||
// }
|
|
||||||
//
|
|
||||||
// /**
|
|
||||||
// * 执行单个流水线节点
|
|
||||||
// */
|
|
||||||
// private void executePipelineNode(String pipelineId, PipelineNodeVo node) {
|
|
||||||
// String serviceName = PIPELINE_SERVICE_PREFIX + node.getBizCode();
|
|
||||||
// SerialPipelineNodeHandler<Object, Object> handler = SpringBeanHolder.getBean(serviceName);
|
|
||||||
// try {
|
|
||||||
// Object result = handler.doProcess(pipelineId, node);
|
|
||||||
// log.info("Pipeline [{}] node [{}] executed, result: {}", pipelineId, node.getName(), result);
|
|
||||||
// } catch (BadRequestException e) {
|
|
||||||
// handlePipelineError(pipelineId, "Business error in pipeline", e);
|
|
||||||
// } catch (Exception e) {
|
|
||||||
// handlePipelineError(pipelineId, "Unexpected error in pipeline", e);
|
|
||||||
// }
|
|
||||||
// }
|
|
||||||
//
|
|
||||||
// /**
|
|
||||||
// * 处理流水线错误
|
|
||||||
// */
|
|
||||||
// private void handlePipelineError(String pipelineId, String message, Exception e) {
|
|
||||||
// log.error(message, e);
|
|
||||||
// if (isRunning(pipelineId)) {
|
|
||||||
// stopForcibly(pipelineId);
|
|
||||||
// }
|
|
||||||
// }
|
|
||||||
//
|
|
||||||
// /**
|
|
||||||
// * 强制停止流水线
|
|
||||||
// */
|
|
||||||
// public void stopForcibly(String pipelineId) {
|
|
||||||
// Optional.ofNullable(stopFlags.get(pipelineId))
|
|
||||||
// .ifPresent(flag -> flag.set(true));
|
|
||||||
//
|
|
||||||
// Optional.ofNullable(runningTasks.get(pipelineId))
|
|
||||||
// .ifPresent(thread -> {
|
|
||||||
// try {
|
|
||||||
// thread.interrupt();
|
|
||||||
// } catch (SecurityException e) {
|
|
||||||
// log.error("Failed to interrupt pipeline [{}]", pipelineId, e);
|
|
||||||
// }
|
|
||||||
// });
|
|
||||||
// }
|
|
||||||
//
|
|
||||||
// /**
|
|
||||||
// * 初始化测试流水线节点
|
|
||||||
// */
|
|
||||||
// public List<PipelineNodeVo> initialize() {
|
|
||||||
// List<PipelineNodeVo> nodes = new ArrayList<>();
|
|
||||||
//
|
|
||||||
// nodes.add(createNode("初始化", "create_release_branch"));
|
|
||||||
// nodes.add(createNode("代码构建", "build_code"));
|
|
||||||
//
|
|
||||||
// return nodes;
|
|
||||||
// }
|
|
||||||
//
|
|
||||||
// /**
|
|
||||||
// * 创建流水线节点
|
|
||||||
// */
|
|
||||||
// private PipelineNodeVo createNode(String name, String bizCode) {
|
|
||||||
// PipelineNodeVo node = new PipelineNodeVo();
|
|
||||||
// node.setName(name);
|
|
||||||
// node.setClick(true);
|
|
||||||
// node.setButtonList(null);
|
|
||||||
// node.setBizCode(bizCode);
|
|
||||||
// node.setStartTimeMillis(DateTime.now().getTime());
|
|
||||||
// node.setDurationMillis(0L);
|
|
||||||
// node.setExecuteStatus(RUNNING_STATUS);
|
|
||||||
// node.setResultDesc(null);
|
|
||||||
// return node;
|
|
||||||
// }
|
|
||||||
//}
|
|
||||||
-245
@@ -1,245 +0,0 @@
|
|||||||
//package cn.odboy.system.framework.flow.core;
|
|
||||||
//
|
|
||||||
//import java.util.Map;
|
|
||||||
//import java.util.concurrent.*;
|
|
||||||
//import java.util.concurrent.atomic.AtomicBoolean;
|
|
||||||
//import java.util.concurrent.atomic.AtomicInteger;
|
|
||||||
//
|
|
||||||
///**
|
|
||||||
// * 流水线任务池
|
|
||||||
// *
|
|
||||||
// * @author odboy
|
|
||||||
// */
|
|
||||||
//public class PipelineTaskPool {
|
|
||||||
// // 默认线程池参数
|
|
||||||
// private static final int DEFAULT_CORE_POOL_SIZE = 5;
|
|
||||||
// private static final int DEFAULT_MAX_POOL_SIZE = 10;
|
|
||||||
// private static final long DEFAULT_KEEP_ALIVE_TIME = 60L;
|
|
||||||
// private static final TimeUnit DEFAULT_TIME_UNIT = TimeUnit.SECONDS;
|
|
||||||
// private static final int DEFAULT_QUEUE_CAPACITY = 100;
|
|
||||||
//
|
|
||||||
// // 线程池实例
|
|
||||||
// private final ThreadPoolExecutor executor;
|
|
||||||
// // 存储正在运行的任务
|
|
||||||
// private final Map<String, ManagedTask> runningTasks = new ConcurrentHashMap<>();
|
|
||||||
// // 线程池管理器名称
|
|
||||||
// private final String name;
|
|
||||||
//
|
|
||||||
// /**
|
|
||||||
// * 使用默认参数创建线程池管理器
|
|
||||||
// */
|
|
||||||
// public PipelineTaskPool(String name) {
|
|
||||||
// this(name, DEFAULT_CORE_POOL_SIZE, DEFAULT_MAX_POOL_SIZE,
|
|
||||||
// DEFAULT_KEEP_ALIVE_TIME, DEFAULT_TIME_UNIT,
|
|
||||||
// new LinkedBlockingQueue<>(DEFAULT_QUEUE_CAPACITY));
|
|
||||||
// }
|
|
||||||
//
|
|
||||||
// /**
|
|
||||||
// * 自定义参数创建线程池管理器
|
|
||||||
// */
|
|
||||||
// public PipelineTaskPool(String name, int corePoolSize, int maxPoolSize,
|
|
||||||
// long keepAliveTime, TimeUnit unit,
|
|
||||||
// BlockingQueue<Runnable> workQueue) {
|
|
||||||
// this.name = name;
|
|
||||||
// this.executor = new ThreadPoolExecutor(
|
|
||||||
// corePoolSize,
|
|
||||||
// maxPoolSize,
|
|
||||||
// keepAliveTime,
|
|
||||||
// unit,
|
|
||||||
// workQueue,
|
|
||||||
// new ThreadFactory() {
|
|
||||||
// private final AtomicInteger threadNumber = new AtomicInteger(1);
|
|
||||||
//
|
|
||||||
// @Override
|
|
||||||
// public Thread newThread(Runnable r) {
|
|
||||||
// return new Thread(r, name + "-thread-" + threadNumber.getAndIncrement());
|
|
||||||
// }
|
|
||||||
// },
|
|
||||||
// new ThreadPoolExecutor.AbortPolicy()
|
|
||||||
// );
|
|
||||||
// }
|
|
||||||
//
|
|
||||||
// /**
|
|
||||||
// * 提交任务并管理
|
|
||||||
// *
|
|
||||||
// * @param taskId 任务ID
|
|
||||||
// * @param task 任务逻辑
|
|
||||||
// * @return true表示提交成功,false表示任务已存在
|
|
||||||
// */
|
|
||||||
// public boolean submitTask(String taskId, Runnable task) {
|
|
||||||
// if (runningTasks.containsKey(taskId)) {
|
|
||||||
// return false;
|
|
||||||
// }
|
|
||||||
//
|
|
||||||
// ManagedTask managedTask = new ManagedTask(taskId, task);
|
|
||||||
// runningTasks.put(taskId, managedTask);
|
|
||||||
// executor.execute(managedTask);
|
|
||||||
// return true;
|
|
||||||
// }
|
|
||||||
//
|
|
||||||
// /**
|
|
||||||
// * 正常停止任务
|
|
||||||
// *
|
|
||||||
// * @param taskId 任务ID
|
|
||||||
// * @return true表示停止成功,false表示任务不存在或已停止
|
|
||||||
// */
|
|
||||||
// public boolean stopTaskGracefully(String taskId) {
|
|
||||||
// ManagedTask task = runningTasks.get(taskId);
|
|
||||||
// if (task != null) {
|
|
||||||
// return task.stopGracefully();
|
|
||||||
// }
|
|
||||||
// return false;
|
|
||||||
// }
|
|
||||||
//
|
|
||||||
// /**
|
|
||||||
// * 强制停止任务
|
|
||||||
// *
|
|
||||||
// * @param taskId 任务ID
|
|
||||||
// * @return true表示停止成功,false表示任务不存在
|
|
||||||
// */
|
|
||||||
// public boolean stopTaskForcibly(String taskId) {
|
|
||||||
// ManagedTask task = runningTasks.remove(taskId);
|
|
||||||
// if (task != null) {
|
|
||||||
// return task.stopForcibly();
|
|
||||||
// }
|
|
||||||
// return false;
|
|
||||||
// }
|
|
||||||
//
|
|
||||||
// /**
|
|
||||||
// * 检查任务是否正在运行
|
|
||||||
// *
|
|
||||||
// * @param taskId 任务ID
|
|
||||||
// * @return true表示正在运行
|
|
||||||
// */
|
|
||||||
// public boolean isTaskRunning(String taskId) {
|
|
||||||
// ManagedTask task = runningTasks.get(taskId);
|
|
||||||
// return task != null && !task.isStopped();
|
|
||||||
// }
|
|
||||||
//
|
|
||||||
// /**
|
|
||||||
// * 获取正在运行的任务数量
|
|
||||||
// *
|
|
||||||
// * @return 运行中任务数
|
|
||||||
// */
|
|
||||||
// public int getRunningTaskCount() {
|
|
||||||
// return runningTasks.size();
|
|
||||||
// }
|
|
||||||
//
|
|
||||||
// /**
|
|
||||||
// * 关闭线程池(等待所有任务完成)
|
|
||||||
// */
|
|
||||||
// public void shutdown() {
|
|
||||||
// executor.shutdown();
|
|
||||||
// }
|
|
||||||
//
|
|
||||||
// /**
|
|
||||||
// * 立即关闭线程池(尝试停止所有任务)
|
|
||||||
// */
|
|
||||||
// public void shutdownNow() {
|
|
||||||
// executor.shutdownNow();
|
|
||||||
// runningTasks.clear();
|
|
||||||
// }
|
|
||||||
//
|
|
||||||
// /**
|
|
||||||
// * 获取线程池状态信息
|
|
||||||
// *
|
|
||||||
// * @return 状态信息字符串
|
|
||||||
// */
|
|
||||||
// public String getPoolStatus() {
|
|
||||||
// return String.format(
|
|
||||||
// "[%s] PoolStatus: Active=%d, Completed=%d, Task=%d, Queue=%d/%d",
|
|
||||||
// name,
|
|
||||||
// executor.getActiveCount(),
|
|
||||||
// executor.getCompletedTaskCount(),
|
|
||||||
// executor.getTaskCount(),
|
|
||||||
// executor.getQueue().size(),
|
|
||||||
// executor.getQueue().remainingCapacity()
|
|
||||||
// );
|
|
||||||
// }
|
|
||||||
//
|
|
||||||
// /**
|
|
||||||
// * 被管理的任务包装类
|
|
||||||
// */
|
|
||||||
// private class ManagedTask implements Runnable {
|
|
||||||
// private final String taskId;
|
|
||||||
// private final Runnable task;
|
|
||||||
// private final AtomicBoolean running = new AtomicBoolean(true);
|
|
||||||
// private final AtomicBoolean stopped = new AtomicBoolean(false);
|
|
||||||
// private Thread currentThread;
|
|
||||||
//
|
|
||||||
// public ManagedTask(String taskId, Runnable task) {
|
|
||||||
// this.taskId = taskId;
|
|
||||||
// this.task = task;
|
|
||||||
// }
|
|
||||||
//
|
|
||||||
// @Override
|
|
||||||
// public void run() {
|
|
||||||
// currentThread = Thread.currentThread();
|
|
||||||
// try {
|
|
||||||
// // 执行前检查是否已被停止
|
|
||||||
// if (running.get()) {
|
|
||||||
// task.run();
|
|
||||||
// }
|
|
||||||
// } finally {
|
|
||||||
// // 任务完成后从运行列表中移除
|
|
||||||
// runningTasks.remove(taskId);
|
|
||||||
// stopped.set(true);
|
|
||||||
// }
|
|
||||||
// }
|
|
||||||
//
|
|
||||||
// /**
|
|
||||||
// * 正常停止任务
|
|
||||||
// *
|
|
||||||
// * @return 是否成功停止
|
|
||||||
// */
|
|
||||||
// public boolean stopGracefully() {
|
|
||||||
// if (running.compareAndSet(true, false)) {
|
|
||||||
// return true;
|
|
||||||
// }
|
|
||||||
// return false;
|
|
||||||
// }
|
|
||||||
//
|
|
||||||
// /**
|
|
||||||
// * 强制停止任务
|
|
||||||
// *
|
|
||||||
// * @return 是否成功停止
|
|
||||||
// */
|
|
||||||
// public boolean stopForcibly() {
|
|
||||||
// if (stopped.get()) {
|
|
||||||
// return false;
|
|
||||||
// }
|
|
||||||
//
|
|
||||||
// if (running.compareAndSet(true, false)) {
|
|
||||||
// if (currentThread != null) {
|
|
||||||
// // 暴力停止线程
|
|
||||||
// // currentThread.stop();
|
|
||||||
// try {
|
|
||||||
// /**
|
|
||||||
// * Thread.interrupt只支持终止线程的阻塞状态(wait、join、sleep),
|
|
||||||
// * 在阻塞出抛出InterruptedException异常,但是并不会终止运行的线程本身;
|
|
||||||
// * 所以需要注意,此处彻底销毁本线程,需要通过共享变量方式;
|
|
||||||
// */
|
|
||||||
// currentThread.interrupt();
|
|
||||||
// currentThread.join();
|
|
||||||
// } catch (InterruptedException e) {
|
|
||||||
// currentThread.interrupt();
|
|
||||||
// } catch (Exception e) {
|
|
||||||
// // ignore
|
|
||||||
// // currentThread.stop();
|
|
||||||
// }
|
|
||||||
// }
|
|
||||||
// return true;
|
|
||||||
// }
|
|
||||||
// return false;
|
|
||||||
// }
|
|
||||||
//
|
|
||||||
// /**
|
|
||||||
// * 检查任务是否已停止
|
|
||||||
// *
|
|
||||||
// * @return true表示已停止
|
|
||||||
// */
|
|
||||||
// public boolean isStopped() {
|
|
||||||
// return stopped.get();
|
|
||||||
// }
|
|
||||||
// }
|
|
||||||
//}
|
|
||||||
-129
@@ -1,129 +0,0 @@
|
|||||||
//package cn.odboy.system.framework.flow.core.handler;
|
|
||||||
//
|
|
||||||
//import cn.hutool.core.thread.ThreadUtil;
|
|
||||||
//import cn.odboy.exception.BadRequestException;
|
|
||||||
//import com.alibaba.fastjson2.JSON;
|
|
||||||
//import lombok.extern.slf4j.Slf4j;
|
|
||||||
//
|
|
||||||
///**
|
|
||||||
// * 串行流水线节点处理器
|
|
||||||
// *
|
|
||||||
// * @author odboy
|
|
||||||
// */
|
|
||||||
//@Slf4j
|
|
||||||
//public abstract class SerialPipelineNodeHandler<I, T> {
|
|
||||||
//
|
|
||||||
// /**
|
|
||||||
// * 预处理
|
|
||||||
// *
|
|
||||||
// * @param pipelineId 流水线id
|
|
||||||
// * @param inputModel 输入参数
|
|
||||||
// */
|
|
||||||
// protected abstract void preProcess(String pipelineId, I inputModel);
|
|
||||||
//
|
|
||||||
// /**
|
|
||||||
// * 预处理异常
|
|
||||||
// *
|
|
||||||
// * @param pipelineId 流水线id
|
|
||||||
// * @param inputModel 输入参数
|
|
||||||
// * @param e 异常信息
|
|
||||||
// */
|
|
||||||
// protected abstract void preProcessException(String pipelineId, I inputModel, Exception e);
|
|
||||||
//
|
|
||||||
// /**
|
|
||||||
// * 处理业务
|
|
||||||
// *
|
|
||||||
// * @param pipelineId 流水线id
|
|
||||||
// * @param inputModel 输入参数
|
|
||||||
// */
|
|
||||||
// protected abstract T process(String pipelineId, I inputModel);
|
|
||||||
//
|
|
||||||
// /**
|
|
||||||
// * 处理业务异常
|
|
||||||
// *
|
|
||||||
// * @param pipelineId 流水线id
|
|
||||||
// * @param inputModel 输入参数
|
|
||||||
// * @param e 异常信息
|
|
||||||
// */
|
|
||||||
// protected abstract void processException(String pipelineId, I inputModel, Exception e);
|
|
||||||
//
|
|
||||||
// /**
|
|
||||||
// * 业务后置处理
|
|
||||||
// *
|
|
||||||
// * @param pipelineId 流水线id
|
|
||||||
// * @param inputModel 输入参数
|
|
||||||
// */
|
|
||||||
// protected abstract void postProcess(String pipelineId, I inputModel, T processResult);
|
|
||||||
//
|
|
||||||
// /**
|
|
||||||
// * 业务后置处理异常
|
|
||||||
// *
|
|
||||||
// * @param pipelineId 流水线id
|
|
||||||
// * @param inputModel 输入参数
|
|
||||||
// */
|
|
||||||
// protected abstract void postProcessException(String pipelineId, I inputModel, T processResult, Exception e);
|
|
||||||
//
|
|
||||||
// /**
|
|
||||||
// * 当前节点是否被锁住(卡点功能),返回false解锁
|
|
||||||
// *
|
|
||||||
// * @param pipelineId 流水线id
|
|
||||||
// * @param inputModel 输入参数
|
|
||||||
// */
|
|
||||||
// protected abstract boolean isLocked(String pipelineId, I inputModel);
|
|
||||||
//
|
|
||||||
// /**
|
|
||||||
// * 当进程被中断时
|
|
||||||
// *
|
|
||||||
// * @param pipelineId 流水线id
|
|
||||||
// * @param inputModel 输入参数
|
|
||||||
// */
|
|
||||||
// protected abstract void onProcessInterrupted(String pipelineId, I inputModel);
|
|
||||||
//
|
|
||||||
// /**
|
|
||||||
// * 业务处理
|
|
||||||
// *
|
|
||||||
// * @param pipelineId 流水线id
|
|
||||||
// * @param inputModel 输入参数
|
|
||||||
// * @return T
|
|
||||||
// */
|
|
||||||
// public T doProcess(String pipelineId, I inputModel) throws BadRequestException {
|
|
||||||
// try {
|
|
||||||
// if (Thread.currentThread().isInterrupted()) {
|
|
||||||
// onProcessInterrupted(pipelineId, inputModel);
|
|
||||||
// throw new InterruptedException("业务被中断");
|
|
||||||
// }
|
|
||||||
// preProcess(pipelineId, inputModel);
|
|
||||||
// log.info("preProcess,pipelineId={},inputModel={}", pipelineId, JSON.toJSONString(inputModel));
|
|
||||||
// } catch (Exception e) {
|
|
||||||
// preProcessException(pipelineId, inputModel, e);
|
|
||||||
// log.error("preProcessException,pipelineId={},inputModel={}", pipelineId, JSON.toJSONString(inputModel), e);
|
|
||||||
// throw new BadRequestException(e.getMessage());
|
|
||||||
// }
|
|
||||||
// if (Thread.currentThread().isInterrupted()) {
|
|
||||||
// throw new BadRequestException("业务被中断");
|
|
||||||
// }
|
|
||||||
// while (isLocked(pipelineId, inputModel)) {
|
|
||||||
// ThreadUtil.safeSleep(2000);
|
|
||||||
// }
|
|
||||||
// try {
|
|
||||||
// if (Thread.currentThread().isInterrupted()) {
|
|
||||||
// throw new BadRequestException("业务被中断");
|
|
||||||
// }
|
|
||||||
// T processResult = process(pipelineId, inputModel);
|
|
||||||
// log.info("process,pipelineId={},inputModel={}", pipelineId, JSON.toJSONString(inputModel));
|
|
||||||
// try {
|
|
||||||
// postProcess(pipelineId, inputModel, processResult);
|
|
||||||
// log.info("postProcess,pipelineId={},inputModel={},processResult={}", pipelineId, JSON.toJSONString(inputModel), JSON.toJSONString(processResult));
|
|
||||||
// return processResult;
|
|
||||||
// } catch (Exception e) {
|
|
||||||
// postProcessException(pipelineId, inputModel, processResult, e);
|
|
||||||
// log.error("postProcessException,pipelineId={},inputModel={},processResult={}", pipelineId, JSON.toJSONString(inputModel), JSON.toJSONString(processResult));
|
|
||||||
// throw new BadRequestException(e.getMessage());
|
|
||||||
// }
|
|
||||||
// } catch (Exception e) {
|
|
||||||
// processException(pipelineId, inputModel, e);
|
|
||||||
// log.error("processException,pipelineId={},inputModel={}", pipelineId, JSON.toJSONString(inputModel), e);
|
|
||||||
// throw new BadRequestException(e.getMessage());
|
|
||||||
// }
|
|
||||||
// }
|
|
||||||
//}
|
|
||||||
-44
@@ -1,44 +0,0 @@
|
|||||||
//package cn.odboy.system.framework.flow.core.handler.impl;
|
|
||||||
//
|
|
||||||
//import cn.odboy.system.framework.flow.core.handler.SerialPipelineNodeHandler;
|
|
||||||
//import lombok.extern.slf4j.Slf4j;
|
|
||||||
//import org.springframework.stereotype.Service;
|
|
||||||
//
|
|
||||||
//@Slf4j
|
|
||||||
//@Service("pipeline:create_release_branch")
|
|
||||||
//public class PipelineAppCreateReleaseBranchHandlerServiceImpl extends SerialPipelineNodeHandler<Object, Object> {
|
|
||||||
// @Override
|
|
||||||
// public void preProcess(String pipelineId, Object inputModel) {
|
|
||||||
// }
|
|
||||||
//
|
|
||||||
// @Override
|
|
||||||
// public void preProcessException(String pipelineId, Object inputModel, Exception e) {
|
|
||||||
// }
|
|
||||||
//
|
|
||||||
// @Override
|
|
||||||
// public Object process(String pipelineId, Object inputModel) {
|
|
||||||
// return "Hello World";
|
|
||||||
// }
|
|
||||||
//
|
|
||||||
// @Override
|
|
||||||
// public void processException(String pipelineId, Object inputModel, Exception e) {
|
|
||||||
// }
|
|
||||||
//
|
|
||||||
// @Override
|
|
||||||
// public void postProcess(String pipelineId, Object inputModel, Object processResult) {
|
|
||||||
// }
|
|
||||||
//
|
|
||||||
// @Override
|
|
||||||
// public void postProcessException(String pipelineId, Object inputModel, Object processResult, Exception e) {
|
|
||||||
// }
|
|
||||||
//
|
|
||||||
// @Override
|
|
||||||
// public boolean isLocked(String pipelineId, Object inputModel) {
|
|
||||||
// return false;
|
|
||||||
// }
|
|
||||||
//
|
|
||||||
// @Override
|
|
||||||
// protected void onProcessInterrupted(String pipelineId, Object inputModel) {
|
|
||||||
// log.error("pipeline is stopping");
|
|
||||||
// }
|
|
||||||
//}
|
|
||||||
-68
@@ -1,68 +0,0 @@
|
|||||||
//package cn.odboy.system.framework.flow.model;
|
|
||||||
//
|
|
||||||
//import cn.odboy.base.CsSerializeObject;
|
|
||||||
//import lombok.Data;
|
|
||||||
//import lombok.EqualsAndHashCode;
|
|
||||||
//
|
|
||||||
//import java.util.List;
|
|
||||||
//
|
|
||||||
//@Data
|
|
||||||
//@EqualsAndHashCode(callSuper = false)
|
|
||||||
//public class PipelineNodeVo extends CsSerializeObject {
|
|
||||||
// /**
|
|
||||||
// * 业务节点名称
|
|
||||||
// */
|
|
||||||
// private String name;
|
|
||||||
// /**
|
|
||||||
// * 是否可以点击
|
|
||||||
// */
|
|
||||||
// private Boolean click = false;
|
|
||||||
// /**
|
|
||||||
// * 功能按钮列表,与click互斥
|
|
||||||
// */
|
|
||||||
// private List<PipelineNodeOperationButton> buttonList;
|
|
||||||
// /**
|
|
||||||
// * 业务编码(createReleaseBranch、mergeCode)
|
|
||||||
// */
|
|
||||||
// private String bizCode;
|
|
||||||
// /**
|
|
||||||
// * 业务节点开始时间
|
|
||||||
// */
|
|
||||||
// private Long startTimeMillis;
|
|
||||||
// /**
|
|
||||||
// * 业务节点耗时
|
|
||||||
// */
|
|
||||||
// private Long durationMillis;
|
|
||||||
// /**
|
|
||||||
// * 业务节点执行状态(success 成功, error 失败, running 运行中)
|
|
||||||
// */
|
|
||||||
// private String executeStatus;
|
|
||||||
// /**
|
|
||||||
// * 业务节点结果描述(执行成功)
|
|
||||||
// */
|
|
||||||
// private String resultDesc;
|
|
||||||
//
|
|
||||||
// /**
|
|
||||||
// * 流水线节点功能按钮
|
|
||||||
// */
|
|
||||||
// @Data
|
|
||||||
// @EqualsAndHashCode(callSuper = false)
|
|
||||||
// public static class PipelineNodeOperationButton extends CsSerializeObject {
|
|
||||||
// /**
|
|
||||||
// * 请求方式(get、post)
|
|
||||||
// */
|
|
||||||
// private String method;
|
|
||||||
// /**
|
|
||||||
// * 请求路径(/api/doCheck?id=1&bizCode=)
|
|
||||||
// */
|
|
||||||
// private String requestUrl;
|
|
||||||
// /**
|
|
||||||
// * 按钮名称(通过、打回)
|
|
||||||
// */
|
|
||||||
// private String text;
|
|
||||||
// /**
|
|
||||||
// * 按钮类型(execute 请求url、link 跳转三方)
|
|
||||||
// */
|
|
||||||
// private String type;
|
|
||||||
// }
|
|
||||||
//}
|
|
||||||
+1
-1
@@ -6,7 +6,7 @@ import org.springframework.web.socket.server.standard.ServerEndpointExporter;
|
|||||||
|
|
||||||
|
|
||||||
@Configuration
|
@Configuration
|
||||||
public class WebSocketConfig {
|
public class WsConfig {
|
||||||
|
|
||||||
@Bean
|
@Bean
|
||||||
public ServerEndpointExporter serverEndpointExporter() {
|
public ServerEndpointExporter serverEndpointExporter() {
|
||||||
+1
-1
@@ -13,7 +13,7 @@ import org.springframework.stereotype.Component;
|
|||||||
* @date 2024-06-07
|
* @date 2024-06-07
|
||||||
*/
|
*/
|
||||||
@Component
|
@Component
|
||||||
public class CsWebSocketFactory implements WebServerFactoryCustomizer<UndertowServletWebServerFactory> {
|
public class WsFactory implements WebServerFactoryCustomizer<UndertowServletWebServerFactory> {
|
||||||
|
|
||||||
@Override
|
@Override
|
||||||
public void customize(UndertowServletWebServerFactory factory) {
|
public void customize(UndertowServletWebServerFactory factory) {
|
||||||
+1
-1
@@ -1,7 +1,7 @@
|
|||||||
package cn.odboy.framework.websocket.constant;
|
package cn.odboy.framework.websocket.constant;
|
||||||
|
|
||||||
|
|
||||||
public enum CsWebSocketBizTypeEnum {
|
public enum CsWsBizTypeEnum {
|
||||||
/**
|
/**
|
||||||
* 连接
|
* 连接
|
||||||
*/
|
*/
|
||||||
-30
@@ -1,30 +0,0 @@
|
|||||||
package cn.odboy.framework.websocket.context;
|
|
||||||
|
|
||||||
import java.util.Collection;
|
|
||||||
import java.util.Map;
|
|
||||||
import java.util.concurrent.ConcurrentHashMap;
|
|
||||||
|
|
||||||
public class CsWebSocketClientManager {
|
|
||||||
/**
|
|
||||||
* concurrent包的线程安全Map,用来存放每个客户端对应的MyWebSocket对象。(分布式必出问题)
|
|
||||||
*/
|
|
||||||
private static final Map<String, WebSocketServer> CLIENT = new ConcurrentHashMap<>();
|
|
||||||
|
|
||||||
public static void addClient(String sid, WebSocketServer webSocketServer) {
|
|
||||||
// 如果存在就先删除一个,防止重复推送消息
|
|
||||||
CLIENT.remove(sid);
|
|
||||||
CLIENT.put(sid, webSocketServer);
|
|
||||||
}
|
|
||||||
|
|
||||||
public static void removeClient(String sid) {
|
|
||||||
CLIENT.remove(sid);
|
|
||||||
}
|
|
||||||
|
|
||||||
public static Collection<WebSocketServer> getAllClient() {
|
|
||||||
return CLIENT.values();
|
|
||||||
}
|
|
||||||
|
|
||||||
public static WebSocketServer getClientBySid(String sid) {
|
|
||||||
return CLIENT.get(sid);
|
|
||||||
}
|
|
||||||
}
|
|
||||||
-15
@@ -1,15 +0,0 @@
|
|||||||
package cn.odboy.framework.websocket.context;
|
|
||||||
|
|
||||||
import cn.odboy.framework.websocket.constant.CsWebSocketBizTypeEnum;
|
|
||||||
import lombok.Data;
|
|
||||||
|
|
||||||
@Data
|
|
||||||
public class CsWebSocketMessage {
|
|
||||||
private String message;
|
|
||||||
private CsWebSocketBizTypeEnum bizType;
|
|
||||||
|
|
||||||
public CsWebSocketMessage(String message, CsWebSocketBizTypeEnum bizType) {
|
|
||||||
this.message = message;
|
|
||||||
this.bizType = bizType;
|
|
||||||
}
|
|
||||||
}
|
|
||||||
+15
@@ -0,0 +1,15 @@
|
|||||||
|
package cn.odboy.framework.websocket.context;
|
||||||
|
|
||||||
|
import cn.odboy.framework.websocket.constant.CsWsBizTypeEnum;
|
||||||
|
import lombok.Data;
|
||||||
|
|
||||||
|
@Data
|
||||||
|
public class CsWsMessage {
|
||||||
|
private String message;
|
||||||
|
private CsWsBizTypeEnum bizType;
|
||||||
|
|
||||||
|
public CsWsMessage(String message, CsWsBizTypeEnum bizType) {
|
||||||
|
this.message = message;
|
||||||
|
this.bizType = bizType;
|
||||||
|
}
|
||||||
|
}
|
||||||
+13
-11
@@ -2,6 +2,7 @@ package cn.odboy.framework.websocket.context;
|
|||||||
|
|
||||||
import com.alibaba.fastjson2.JSON;
|
import com.alibaba.fastjson2.JSON;
|
||||||
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.*;
|
||||||
@@ -14,7 +15,9 @@ import java.util.Objects;
|
|||||||
@Slf4j
|
@Slf4j
|
||||||
@Component
|
@Component
|
||||||
@ServerEndpoint("/webSocket/{sid}")
|
@ServerEndpoint("/webSocket/{sid}")
|
||||||
public class WebSocketServer {
|
public class CsWsServer {
|
||||||
|
@Autowired
|
||||||
|
private WsClientManager wsClientManager;
|
||||||
/**
|
/**
|
||||||
* 与某个客户端的连接会话,需要通过它来给客户端发送数据
|
* 与某个客户端的连接会话,需要通过它来给客户端发送数据
|
||||||
*/
|
*/
|
||||||
@@ -31,7 +34,7 @@ public class WebSocketServer {
|
|||||||
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;
|
||||||
CsWebSocketClientManager.addClient(sid, this);
|
wsClientManager.addClient(sid, this);
|
||||||
}
|
}
|
||||||
|
|
||||||
/**
|
/**
|
||||||
@@ -39,7 +42,7 @@ public class WebSocketServer {
|
|||||||
*/
|
*/
|
||||||
@OnClose
|
@OnClose
|
||||||
public void onClose() {
|
public void onClose() {
|
||||||
CsWebSocketClientManager.removeClient(this.sid);
|
wsClientManager.removeClient(this.sid);
|
||||||
}
|
}
|
||||||
|
|
||||||
/**
|
/**
|
||||||
@@ -58,8 +61,8 @@ public class WebSocketServer {
|
|||||||
*
|
*
|
||||||
* @param message /
|
* @param message /
|
||||||
*/
|
*/
|
||||||
private static void sendToAll(String message) {
|
private void sendToAll(String message) {
|
||||||
for (WebSocketServer item : CsWebSocketClientManager.getAllClient()) {
|
for (CsWsServer item : wsClientManager.getAllClient()) {
|
||||||
try {
|
try {
|
||||||
item.innerSendMessage(message);
|
item.innerSendMessage(message);
|
||||||
} catch (IOException e) {
|
} catch (IOException e) {
|
||||||
@@ -71,7 +74,6 @@ public class WebSocketServer {
|
|||||||
@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);
|
||||||
error.printStackTrace();
|
|
||||||
}
|
}
|
||||||
|
|
||||||
/**
|
/**
|
||||||
@@ -84,14 +86,14 @@ public class WebSocketServer {
|
|||||||
/**
|
/**
|
||||||
* 群发自定义消息
|
* 群发自定义消息
|
||||||
*/
|
*/
|
||||||
public static void sendMessage(CsWebSocketMessage webSocketMessage, @PathParam("sid") String sid) throws IOException {
|
public void sendMessage(CsWsMessage message, @PathParam("sid") String sid) throws IOException {
|
||||||
String message = JSON.toJSONString(webSocketMessage);
|
String body = JSON.toJSONString(message);
|
||||||
log.info("推送消息到{},推送内容:{}", sid, message);
|
log.info("推送消息到{},推送内容:{}", sid, message);
|
||||||
try {
|
try {
|
||||||
if (sid == null) {
|
if (sid == null) {
|
||||||
sendToAll(message);
|
sendToAll(body);
|
||||||
} else {
|
} else {
|
||||||
CsWebSocketClientManager.getClientBySid(sid).innerSendMessage(message);
|
wsClientManager.getClientBySid(sid).innerSendMessage(body);
|
||||||
}
|
}
|
||||||
} catch (Exception ignored) {
|
} catch (Exception ignored) {
|
||||||
}
|
}
|
||||||
@@ -105,7 +107,7 @@ public class WebSocketServer {
|
|||||||
if (o == null || getClass() != o.getClass()) {
|
if (o == null || getClass() != o.getClass()) {
|
||||||
return false;
|
return false;
|
||||||
}
|
}
|
||||||
WebSocketServer that = (WebSocketServer) 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);
|
||||||
}
|
}
|
||||||
+33
@@ -0,0 +1,33 @@
|
|||||||
|
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);
|
||||||
|
}
|
||||||
|
}
|
||||||
Reference in New Issue
Block a user