From dbea5a7b3fcc2430ffbfc1d64203d76bbf85321e Mon Sep 17 00:00:00 2001 From: tangw Date: Thu, 27 Nov 2025 19:07:55 +0800 Subject: [PATCH] =?UTF-8?q?websocket=E6=B6=88=E6=81=AF=E6=9C=8D=E5=8A=A1?= =?UTF-8?q?=E5=99=A8=E6=A8=A1=E5=9D=97?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../websocket/ChatWebSocketServer.java | 295 ++++++++++++++++++ .../websocket/RetryMsgService.java | 107 +++++++ .../websocket/RetryMsgTask.java | 21 ++ .../websocket/UserService.java | 60 ++++ .../websocket/WebSocketMessage.java | 24 ++ .../websocket/WebSocketMsgService.java | 15 + .../websocket/WebSocketMsgServiceImpl.java | 172 ++++++++++ .../websocket/WebSocketUser.java | 78 +++++ .../czlis/interfaceCommon/websocket/notsbook | 57 ++++ 9 files changed, 829 insertions(+) create mode 100644 lis-interfaceCommon/src/main/java/com/czlis/interfaceCommon/websocket/ChatWebSocketServer.java create mode 100644 lis-interfaceCommon/src/main/java/com/czlis/interfaceCommon/websocket/RetryMsgService.java create mode 100644 lis-interfaceCommon/src/main/java/com/czlis/interfaceCommon/websocket/RetryMsgTask.java create mode 100644 lis-interfaceCommon/src/main/java/com/czlis/interfaceCommon/websocket/UserService.java create mode 100644 lis-interfaceCommon/src/main/java/com/czlis/interfaceCommon/websocket/WebSocketMessage.java create mode 100644 lis-interfaceCommon/src/main/java/com/czlis/interfaceCommon/websocket/WebSocketMsgService.java create mode 100644 lis-interfaceCommon/src/main/java/com/czlis/interfaceCommon/websocket/WebSocketMsgServiceImpl.java create mode 100644 lis-interfaceCommon/src/main/java/com/czlis/interfaceCommon/websocket/WebSocketUser.java create mode 100644 lis-interfaceCommon/src/main/java/com/czlis/interfaceCommon/websocket/notsbook diff --git a/lis-interfaceCommon/src/main/java/com/czlis/interfaceCommon/websocket/ChatWebSocketServer.java b/lis-interfaceCommon/src/main/java/com/czlis/interfaceCommon/websocket/ChatWebSocketServer.java new file mode 100644 index 0000000..832cf6c --- /dev/null +++ b/lis-interfaceCommon/src/main/java/com/czlis/interfaceCommon/websocket/ChatWebSocketServer.java @@ -0,0 +1,295 @@ +package com.czlis.interfaceCommon.websocket; + +import com.alibaba.fastjson2.JSON; +import com.czlis.common.core.domain.model.LoginUser; +import com.czlis.common.utils.SecurityUtils; +import com.czlis.common.utils.spring.SpringUtils; +import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.web.socket.server.standard.SpringConfigurator; + +import javax.websocket.*; +import javax.websocket.server.PathParam; +import javax.websocket.server.ServerEndpoint; +import java.net.InetSocketAddress; +import java.util.*; +import java.util.concurrent.ConcurrentHashMap; + +/** + * lis消息中心服务器 + * 原理:核心是 统一管理所有在线 WebSocket 连接(而非直接引用实例)—— 因为 @ServerEndpoint 是「多实例」模式(每个客户端连接对应一个独立实例),直接引用单个实例毫无意义,且会导致重复创建 / 连接泄漏问题。 + * 业务结构:在 WebSocket 端点类内部维护一个「全局共享的在线连接容器」,其他业务通过这个容器获取目标连接并发送消息 + * 一、核心设计原则 + * 全局连接容器:用线程安全的集合(如 ConcurrentHashMap)存储所有在线 Session(WebSocket 连接会话,每个客户端唯一),避免多线程并发问题; + * 实例无状态:@ServerEndpoint 实例仅用于处理连接生命周期,不存储业务数据,所有共享状态(如在线连接)放在「静态容器」中; + * 工具方法暴露:在 WebSocket 端点类中提供静态工具方法(如 sendToUser、broadcast),其他业务直接调用该方法,无需关心实例创建; + * 避免重复连接:连接的创建 / 销毁由 WebSocket 协议自动管理(@OnOpen 新增、@OnClose 移除),业务层不干预实例创建。 + */ +@ServerEndpoint(value = "/ws/msgserver/{userId}", configurator = SpringConfigurator.class) +public class ChatWebSocketServer { + // 🌟 核心固定参数白名单:这些参数不存入动态属性(单独解析为固定身份属性) + private static final Set CORE_PARAMS = new HashSet<>(Arrays.asList("token", "username", "role")); + // 🌟 容器改为存储封装后的WebSocketUser(key=userId) + private static final Map ONLINE_USER_MAP = new ConcurrentHashMap<>(); + /** + * 连接建立时:封装用户信息并存入容器 + * 注意:若需要传递username/role等属性,可通过路径参数/请求参数/Token解析 + */ + @OnOpen + public void onOpen(@PathParam("userId") String userId, Session session) { + // 1. 获取所有请求参数(客户端连接时携带的参数,如 ws://xxx/ws/chat/1001?token=xxx&username=张三&role=USER&deviceType=pc&version=1.0) + Map> requestParams = session.getRequestParameterMap(); + String username = ""; + String role = ""; + String loginIp = getClientIp(session); + String token=""; + try { + // 场景1:正常业务场景(前端登录后连接)→ 从 SecurityUtils 获取 + LoginUser loginUser = SecurityUtils.getLoginUser(); + token=loginUser.getToken(); + username = loginUser.getUser().getNickName(); + role = loginUser.getUser().getRoleIds().toString(); + } catch (Exception e) { + // 场景2:测试工具场景(无登录上下文)→ 从 URL 参数解析兜底 + // 从参数中获取 username/role(测试工具连接时手动拼接) + username = requestParams.getOrDefault("username", Collections.singletonList("测试用户")).get(0); + role = requestParams.getOrDefault("role", Collections.singletonList("TEST")).get(0); + System.out.println("未获取到登录上下文,使用URL参数兜底:用户名=" + username + ",角色=" + role); + } + // 3. 封装核心用户信息(固定属性) + WebSocketUser webSocketUser = new WebSocketUser(userId, username, role, loginIp, session); + webSocketUser.setDynamicAttr("token", token); + // 4. 🌟 自动解析剩余参数,存入dynamicAttrs(过滤核心参数) + for (Map.Entry> entry : requestParams.entrySet()) { + String paramKey = entry.getKey(); + // 跳过核心固定参数、空值参数 + if (CORE_PARAMS.contains(paramKey) || entry.getValue() == null || entry.getValue().isEmpty()) { + continue; + } + // 参数值取第一个(WebSocket请求参数通常是单值,多值可自行处理) + String paramValue = entry.getValue().get(0); + // 4. 跳过空字符串/无效值(可选,进一步过滤垃圾数据) + if (paramValue.isEmpty() || "null".equalsIgnoreCase(paramValue)) { + continue; + } + // 自动存入动态属性 + webSocketUser.setDynamicAttr(paramKey, paramValue); + } + // 5. 存入全局容器 + ONLINE_USER_MAP.put(userId, webSocketUser); + System.out.println("用户[" + userId + "-" + username + "]连接成功,IP:" + loginIp + ",当前在线人数:" + ONLINE_USER_MAP.size()); + } + // ========== 消息接收核心方法 ========== + + /** + * 业务层处理客户端通过 WebSocket 发送的消息,核心思路是:将「消息接收」与「业务处理」解耦 —— WebSocket 端点类仅负责接收消息、做基础校验,再将消息转发给业务层处理(同步 / 异步),业务层专注于业务逻辑(如消息解析、权限校验、数据处理、结果推送),最终通过 WebSocket 端点类将处理结果返回给客户端。 + * 整体流程设计 + * 客户端发送消息 → WebSocket端点(@OnMessage)接收 → 基础校验 → 转发业务层处理 → 业务层执行业务逻辑 → 结果通过WebSocket推送给客户端 + * + * @param message + * @param session + * @param userId + */ + @OnMessage + public void onMessage(String message, Session session, @PathParam("userId") String userId) { + try { + // 1. 基础校验:消息为空/连接关闭 → 直接返回 + if (message == null || message.trim().isEmpty() || !session.isOpen()) { + return; + } + + // 2. 解析消息为统一协议对象(JSON解析) + WebSocketMessage wsMessage = JSON.parseObject(message, WebSocketMessage.class); + if (wsMessage == null) { + sendToUser(userId, "错误:消息格式不合法,请按JSON格式发送"); + return; + } + //暂时使用默认时间戳,生产环境应该使用客户端必传字段 + if (wsMessage.getTimestamp() == null || wsMessage.getTimestamp() == 0) { + wsMessage.setTimestamp(System.currentTimeMillis()); + } + // 3. 补充发送者ID(从路径参数获取,避免客户端伪造) + wsMessage.setSenderId(userId); + + // 4. 获取Spring管理的业务层Bean,转发消息处理 + WebSocketMsgService msgService = SpringUtils.getBean(WebSocketMsgService.class); + if (msgService == null) { + throw new RuntimeException("获取WebSocketMsgService Bean失败,请检查Spring配置"); + } + // 同步处理(简单场景)/ 异步处理(高并发场景,推荐) + msgService.handleMessage(wsMessage, session); + + } catch (Exception e) { + e.printStackTrace(); + // 异常反馈给客户端 + RetryMsgService retryMsgService = SpringUtils.getBean(RetryMsgService.class); + // 封装错误消息存入重试队列 + WebSocketMessage errorMsg = new WebSocketMessage(); + errorMsg.setSenderId(userId.toString()); + errorMsg.setMsgType("ERROR"); + errorMsg.setData("消息处理失败:" + e.getMessage()); + errorMsg.setTimestamp(System.currentTimeMillis()); + retryMsgService.sendToRetryQueue(errorMsg); + } + } + + /** + * 连接关闭时:移除用户信息 + */ + @OnClose + public void onClose(@PathParam("userId") String userId, Session session) { + ONLINE_USER_MAP.remove(userId); + System.out.println("用户[" + userId + "]断开连接,当前在线人数:" + ONLINE_USER_MAP.size()); + } + + /** + * 通信异常时:移除用户信息 + */ + @OnError + public void onError(@PathParam("userId") String userId, Session session, Throwable throwable) { + throwable.printStackTrace(); + ONLINE_USER_MAP.remove(userId); + } + + // ========== 工具方法:获取客户端IP ========== + private String getClientIp(Session session) { + try { + // 获取客户端远程地址(不同服务器容器可能有差异) + InetSocketAddress remoteAddress = (InetSocketAddress) session.getUserProperties().get("javax.websocket.endpoint.remoteAddress"); + return remoteAddress.getAddress().getHostAddress(); + } catch (Exception e) { + return "未知IP"; + } + } + + // ========== 🌟 暴露属性查询方法(业务层调用) ========== + + /** + * 获取所有在线用户的完整信息 + */ + public static Collection getAllOnlineUsers() { + // 返回快照,避免并发修改 + return Collections.unmodifiableCollection(new ArrayList<>(ONLINE_USER_MAP.values())); + } + + /** + * 根据userId获取用户完整信息 + */ + public static WebSocketUser getOnlineUser(String userId) { + return ONLINE_USER_MAP.get(userId); + } + + /** + * 根据角色筛选在线用户(如仅获取管理员) + */ + public static List getOnlineUsersByRole(String role) { + List roleUsers = new ArrayList<>(); + for (WebSocketUser user : ONLINE_USER_MAP.values()) { + if (role.equals(user.getRole()) && user.isConnected()) { + roleUsers.add(user); + } + } + return roleUsers; + } + + /** + * 获取在线用户的userId列表(兼容原有逻辑) + */ + public static Set getAllOnlineUserIds() { + return Collections.unmodifiableSet(new HashSet<>(ONLINE_USER_MAP.keySet())); + } + + // ========== 消息发送方法(适配新容器) ========== + public static boolean sendToUser(String userId, String message) { + WebSocketUser user = ONLINE_USER_MAP.get(userId); + if (user == null || !user.isConnected()) { + System.out.println("用户[" + userId + "]不在线或连接已关闭"); + return false; + } + + try { + user.getSession().getAsyncRemote().sendText(message); + return true; + } catch (Exception e) { + e.printStackTrace(); + ONLINE_USER_MAP.remove(userId); + return false; + } + } + + + public static void broadcast(String message) { + for (WebSocketUser user : ONLINE_USER_MAP.values()) { + if (user.isConnected()) { + user.getSession().getAsyncRemote().sendText(message); + } + } + } + // ========== 🌟 封装属性修改方法(仅允许修改动态属性) ========== + /** + * 给指定用户设置单个动态属性(业务层调用) + * + * @param userId 用户ID + * @param attrKey 属性键 + * @param attrValue 属性值 + * @return true=修改成功(用户在线),false=用户不在线 + */ + public static boolean setUserDynamicAttr(String userId, String attrKey, Object attrValue) { + // 1. 校验用户是否在线 + WebSocketUser user = ONLINE_USER_MAP.get(userId); + if (user == null || !user.isConnected()) { + System.out.println("用户[" + userId + "]不在线,无法修改属性"); + return false; + } + + // 2. 禁止修改核心属性(防止业务层误操作) + if (CORE_PARAMS.contains(attrKey)) { + System.out.println("禁止修改核心属性:" + attrKey); + return false; + } + + // 3. 受控修改动态属性(WebSocketUser的dynamicAttrs是线程安全的) + user.setDynamicAttr(attrKey, attrValue); + System.out.println("用户[" + userId + "]的属性[" + attrKey + "]已更新为:" + attrValue); + return true; + } + + /** + * 批量修改指定用户的动态属性 + * + * @param userId 用户ID + * @param attrMap 待修改的属性键值对 + * @return true=修改成功 + */ + public static boolean batchSetUserDynamicAttr(String userId, Map attrMap) { + WebSocketUser user = ONLINE_USER_MAP.get(userId); + if (user == null || !user.isConnected()) { + return false; + } + + // 批量设置(过滤核心属性) + for (Map.Entry entry : attrMap.entrySet()) { + String attrKey = entry.getKey(); + if (!CORE_PARAMS.contains(attrKey)) { + user.setDynamicAttr(attrKey, entry.getValue()); + } + } + return true; + } + + /** + * 移除指定用户的某个动态属性 + * + * @param userId 用户ID + * @param attrKey 属性键 + * @return true=移除成功 + */ + public static boolean removeUserDynamicAttr(String userId, String attrKey) { + WebSocketUser user = ONLINE_USER_MAP.get(userId); + if (user == null || !user.isConnected()) { + return false; + } + user.removeDynamicAttr(attrKey); + return true; + } + +} diff --git a/lis-interfaceCommon/src/main/java/com/czlis/interfaceCommon/websocket/RetryMsgService.java b/lis-interfaceCommon/src/main/java/com/czlis/interfaceCommon/websocket/RetryMsgService.java new file mode 100644 index 0000000..3daf1ec --- /dev/null +++ b/lis-interfaceCommon/src/main/java/com/czlis/interfaceCommon/websocket/RetryMsgService.java @@ -0,0 +1,107 @@ +package com.czlis.interfaceCommon.websocket; + +import com.alibaba.fastjson2.JSON; +import com.czlis.common.core.redis.RedisCache; +import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.stereotype.Service; + +import java.util.concurrent.TimeUnit; + +/** + * 消息重试服务(若依框架适配) + */ +@Service +public class RetryMsgService { + + // Redis键前缀(若依框架建议添加业务前缀,避免键冲突) + private static final String RETRY_QUEUE_KEY = "ws:msg:retry:queue"; + private static final String DEAD_LETTER_KEY = "ws:msg:dead:letter"; + private static final String RETRY_COUNT_KEY_PREFIX = "ws:msg:retry:count:"; + + // 最大重试次数(可配置在若依的application.yml中,用@Value注入) + private static final int MAX_RETRY_COUNT = 3; + @Autowired + private RedisCache redisCache; + @Autowired + private WebSocketMsgService webSocketMsgService; + + /** + * 存入重试队列(使用若依RedisTemplate) + */ + public void sendToRetryQueue(WebSocketMessage wsMessage) { + if (wsMessage == null) { + return; + } + String msgId = generateMsgId(wsMessage); + String msgJson = JSON.toJSONString(wsMessage); + // 若依Redis:存储重试次数 + redisCache.setCacheObject(RETRY_COUNT_KEY_PREFIX + msgId, "0", 24, TimeUnit.HOURS); + // 存入重试队列 + redisCache.redisTemplate.opsForList().leftPush(RETRY_QUEUE_KEY, msgJson); + System.out.println("消息[" + msgId + "]存入重试队列"); + } + + /** + * 定时重试处理(由定时任务调用) + */ + public void retryProcessMsg() { + + // 从队列取出消息(阻塞1秒) + String msgJson = (String) redisCache.redisTemplate.opsForList().rightPop( + RETRY_QUEUE_KEY, + 1, + TimeUnit.SECONDS + ); + if (msgJson == null) { + return; + } + + try { + WebSocketMessage wsMessage = JSON.parseObject(msgJson, WebSocketMessage.class); + String msgId = generateMsgId(wsMessage); + + // 🌟 替换方案2:使用RedisCache封装的getCacheObject方法(更简洁) + String retryCountStr = redisCache.getCacheObject(RETRY_COUNT_KEY_PREFIX + msgId); + int retryCount = retryCountStr == null ? 0 : Integer.parseInt(retryCountStr); + + // 超过最大重试次数,转入死信队列 + if (retryCount >= MAX_RETRY_COUNT) { + redisCache.redisTemplate.opsForList().leftPush(DEAD_LETTER_KEY, msgJson); + redisCache.redisTemplate.delete(RETRY_COUNT_KEY_PREFIX + msgId); + System.out.println("消息[" + msgId + "]重试" + MAX_RETRY_COUNT + "次失败,转入死信队列"); + return; + } + + // 重试处理 + webSocketMsgService.handleMessage(wsMessage, null); + + // 处理成功,删除重试计数 + redisCache.redisTemplate.delete(RETRY_COUNT_KEY_PREFIX + msgId); + System.out.println("消息[" + msgId + "]重试第" + (retryCount + 1) + "次成功"); + } catch (Exception e) { + e.printStackTrace(); + WebSocketMessage wsMessage = JSON.parseObject(msgJson, WebSocketMessage.class); + String msgId = generateMsgId(wsMessage); + + int retryCount = Integer.parseInt(redisCache.getCacheObject(RETRY_COUNT_KEY_PREFIX + msgId)); + retryCount++; + + // 更新重试次数,重新存入队列 + redisCache.setCacheObject( + RETRY_COUNT_KEY_PREFIX + msgId, + String.valueOf(retryCount), + 24, + TimeUnit.HOURS + ); + redisCache.redisTemplate.opsForList().leftPush(RETRY_QUEUE_KEY, msgJson); + System.out.println("消息[" + msgId + "]重试第" + retryCount + "次失败,重新存入队列"); + } + } + + /** + * 生成消息唯一标识 + */ + private String generateMsgId(WebSocketMessage wsMessage) { + return wsMessage.getSenderId() + "_" + wsMessage.getTimestamp(); + } +} \ No newline at end of file diff --git a/lis-interfaceCommon/src/main/java/com/czlis/interfaceCommon/websocket/RetryMsgTask.java b/lis-interfaceCommon/src/main/java/com/czlis/interfaceCommon/websocket/RetryMsgTask.java new file mode 100644 index 0000000..d2c1b60 --- /dev/null +++ b/lis-interfaceCommon/src/main/java/com/czlis/interfaceCommon/websocket/RetryMsgTask.java @@ -0,0 +1,21 @@ +package com.czlis.interfaceCommon.websocket; + +import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.scheduling.annotation.Scheduled; +import org.springframework.stereotype.Component; + +@Component +public class RetryMsgTask { + + @Autowired + private RetryMsgService retryMsgService; + + /** + * 每5秒执行一次(若依框架支持cron表达式或fixedRate) + */ + @Scheduled(fixedRate = 5000) + public void executeRetry() { + // 多线程处理,避免阻塞定时任务线程 + new Thread(retryMsgService::retryProcessMsg).start(); + } +} diff --git a/lis-interfaceCommon/src/main/java/com/czlis/interfaceCommon/websocket/UserService.java b/lis-interfaceCommon/src/main/java/com/czlis/interfaceCommon/websocket/UserService.java new file mode 100644 index 0000000..9db2043 --- /dev/null +++ b/lis-interfaceCommon/src/main/java/com/czlis/interfaceCommon/websocket/UserService.java @@ -0,0 +1,60 @@ +package com.czlis.interfaceCommon.websocket; + +import org.springframework.stereotype.Service; + +import java.util.List; + +@Service +public class UserService { + + /** + * 示例1:获取所有在线管理员并推送通知 + */ + public void pushNoticeToAdmin(String notice) { + // 根据角色筛选在线管理员 + List adminUsers = ChatWebSocketServer.getOnlineUsersByRole("ADMIN"); + if (adminUsers.isEmpty()) { + System.out.println("当前无在线管理员"); + return; + } + + String message = "【管理员通知】" + notice; + for (WebSocketUser admin : adminUsers) { + ChatWebSocketServer.sendToUser(admin.getUserId(), message); + System.out.println("已给管理员[" + admin.getUsername() + "]推送通知,IP:" + admin.getLoginIp()); + } + } + + /** + * 示例2:获取指定用户的详细连接信息 + */ + public void getUserConnectInfo(String userId) { + WebSocketUser user = ChatWebSocketServer.getOnlineUser(userId); + if (user == null) { + System.out.println("用户[" + userId + "]不在线"); + return; + } + + // 获取身份属性 + System.out.println("用户名:" + user.getUsername()); + System.out.println("角色:" + user.getRole()); + System.out.println("登录IP:" + user.getLoginIp()); + System.out.println("连接时间:" + user.getConnectTime()); + System.out.println("是否在线:" + user.isConnected()); + + // 获取动态属性(如设备类型) + String deviceType = (String) user.getDynamicAttr("deviceType"); + System.out.println("登录设备:" + (deviceType == null ? "未知" : deviceType)); + } + + /** + * 示例3:给在线用户设置动态属性(如设备类型) + */ + public void setUserDeviceAttr(String userId, String deviceType) { + WebSocketUser user = ChatWebSocketServer.getOnlineUser(userId); + if (user != null) { + user.setDynamicAttr("deviceType", deviceType); + System.out.println("已给用户[" + userId + "]设置设备类型:" + deviceType); + } + } +} diff --git a/lis-interfaceCommon/src/main/java/com/czlis/interfaceCommon/websocket/WebSocketMessage.java b/lis-interfaceCommon/src/main/java/com/czlis/interfaceCommon/websocket/WebSocketMessage.java new file mode 100644 index 0000000..f9f3f75 --- /dev/null +++ b/lis-interfaceCommon/src/main/java/com/czlis/interfaceCommon/websocket/WebSocketMessage.java @@ -0,0 +1,24 @@ +package com.czlis.interfaceCommon.websocket; + +import lombok.AllArgsConstructor; +import lombok.Data; +import lombok.NoArgsConstructor; + +/** + * 客户端与服务端的统一消息协议 + */ +@Data +@AllArgsConstructor +@NoArgsConstructor +public class WebSocketMessage { + // 消息类型(如CHAT/NOTICE/ORDER/HEARTBEAT) + private String msgType; + // 业务数据(JSON字符串,可解析为对应业务DTO) + private String data; + // 发送者ID(可选,也可从WebSocket连接中获取) + private String senderId; + // 接收者ID(可选,单聊/广播区分) + private String receiverId; + // 消息时间戳 + private Long timestamp; +} \ No newline at end of file diff --git a/lis-interfaceCommon/src/main/java/com/czlis/interfaceCommon/websocket/WebSocketMsgService.java b/lis-interfaceCommon/src/main/java/com/czlis/interfaceCommon/websocket/WebSocketMsgService.java new file mode 100644 index 0000000..5f61717 --- /dev/null +++ b/lis-interfaceCommon/src/main/java/com/czlis/interfaceCommon/websocket/WebSocketMsgService.java @@ -0,0 +1,15 @@ +package com.czlis.interfaceCommon.websocket; + +import javax.websocket.Session; + +/** + * WebSocket消息处理业务接口 + */ +public interface WebSocketMsgService { + /** + * 处理客户端发送的消息 + * @param wsMessage 解析后的统一消息对象 + * @param session 客户端WebSocket会话 + */ + void handleMessage(WebSocketMessage wsMessage, Session session); +} \ No newline at end of file diff --git a/lis-interfaceCommon/src/main/java/com/czlis/interfaceCommon/websocket/WebSocketMsgServiceImpl.java b/lis-interfaceCommon/src/main/java/com/czlis/interfaceCommon/websocket/WebSocketMsgServiceImpl.java new file mode 100644 index 0000000..5c78a65 --- /dev/null +++ b/lis-interfaceCommon/src/main/java/com/czlis/interfaceCommon/websocket/WebSocketMsgServiceImpl.java @@ -0,0 +1,172 @@ +package com.czlis.interfaceCommon.websocket; + +import com.alibaba.fastjson2.JSON; +import com.czlis.common.annotation.Async; +import com.czlis.common.core.redis.RedisCache; +import com.czlis.common.utils.spring.SpringUtils; +import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.context.annotation.Bean; +import org.springframework.scheduling.concurrent.ThreadPoolTaskExecutor; +import org.springframework.stereotype.Service; +import javax.websocket.Session; +import java.util.Map; +import java.util.concurrent.Executor; +import java.util.concurrent.ThreadPoolExecutor; +import java.util.concurrent.TimeUnit; + +/**业务层(Spring Bean)专注于「消息解析、权限校验、业务逻辑处理、结果推送」,支持不同类型消息的差异化处理。 + * WebSocket消息处理业务实现类 + */ +@Service +public class WebSocketMsgServiceImpl implements WebSocketMsgService { + + @Override + public void handleMessage(WebSocketMessage wsMessage, Session session) { + if (wsMessage.getTimestamp() == null || wsMessage.getTimestamp() == 0) { + ChatWebSocketServer.sendToUser(wsMessage.getSenderId(), "错误:消息必须携带时间戳"); + return; + } + String msgId = wsMessage.getSenderId() + "_" + wsMessage.getTimestamp(); + // 校验消息是否已处理(如Redis中是否存在msgId) + if (SpringUtils.getBean(RedisCache.class).redisTemplate.hasKey("ws:msg:processed:" + msgId)) { + ChatWebSocketServer.sendToUser(wsMessage.getSenderId(), "提示:消息已处理,无需重复发送"); + return; + } + try { + // 1. 基础校验:消息类型、发送者ID + String msgType = wsMessage.getMsgType(); + String senderId = wsMessage.getSenderId(); + if (msgType == null || senderId == null) { + ChatWebSocketServer.sendToUser(senderId, "错误:消息类型/发送者ID不能为空"); + return; + } + + // 2. 根据消息类型分发处理(策略模式,可扩展) + switch (msgType) { + case "CHAT": // 单聊消息 + handleChatMsg(wsMessage); + break; + case "BROADCAST": // 广播消息 + handleBroadcastMsg(wsMessage); + break; + case "ORDER": // 订单相关消息 + handleOrderMsg(wsMessage); + break; + case "HEARTBEAT": // 心跳包(仅确认,无需处理业务) + handleHeartbeatMsg(wsMessage); + break; + default: + ChatWebSocketServer.sendToUser(senderId, "错误:不支持的消息类型:" + msgType); + } + } catch (Exception e) { + e.printStackTrace(); + // ChatWebSocketServer.sendToUser(wsMessage.getSenderId(), "错误:消息处理失败:" + e.getMessage()); + String senderId = wsMessage.getSenderId(); + if (senderId != null && session != null && session.isOpen()) { + ChatWebSocketServer.sendToUser(senderId, "错误:消息处理失败,已加入重试队列"); + } + RetryMsgService retryMsgService = SpringUtils.getBean(RetryMsgService.class); + retryMsgService.sendToRetryQueue(wsMessage); + } + // 标记消息已处理(设置过期时间,避免内存泄漏) + SpringUtils.getBean(RedisCache.class).redisTemplate.opsForValue().set("ws:msg:processed:" + msgId, "1", 24, TimeUnit.HOURS); + } + + /** + * 处理单聊消息 + */ + private void handleChatMsg(WebSocketMessage wsMessage) { + /* + String senderId = wsMessage.getSenderId(); + String receiverId = wsMessage.getReceiverId(); + String data = wsMessage.getData(); + + // 1. 校验接收者是否在线 + boolean receiverOnline = ChatWebSocketServer.getOnlineUser(receiverId) != null; + if (!receiverOnline) { + ChatWebSocketServer.sendToUser(senderId, "提示:接收者[" + receiverId + "]不在线,消息已存入离线队列"); + // 调用离线消息服务:存入Redis/MQ,等待接收者上线后推送 + // userService.saveOfflineMessage(receiverId, data); + return; + } + + // 2. 组装转发消息(添加发送者信息) + Map chatData = JSON.parseObject(data); + // chatData.put("senderName", userService.getUsernameById(senderId)); // 获取发送者昵称 + String forwardMsg = JSON.toJSONString(chatData); + + // 3. 转发消息给接收者 + boolean sendSuccess = ChatWebSocketServer.sendToUser(receiverId, forwardMsg); + if (sendSuccess) { + // 给发送者反馈:消息已送达 + ChatWebSocketServer.sendToUser(senderId, "提示:消息已送达接收者[" + receiverId + "]"); + } + + */ + } + + /** + * 处理广播消息(仅管理员可发送) + */ + private void handleBroadcastMsg(WebSocketMessage wsMessage) { + /* + String senderId = wsMessage.getSenderId(); + String data = wsMessage.getData(); + + // 1. 权限校验:仅管理员可发送广播 + boolean isAdmin = "ADMIN".equals(userService.getUserRoleById(senderId)); + if (!isAdmin) { + ChatWebSocketServer.sendToUser(senderId, "错误:无广播消息权限"); + return; + } + + // 2. 广播消息给所有在线用户 + String broadcastMsg = "【系统广播】" + data; + ChatWebSocketServer.broadcast(broadcastMsg); + + // 3. 给发送者反馈 + ChatWebSocketServer.sendToUser(senderId, "提示:广播消息已发送,覆盖在线人数:" + ChatWebSocketServer.getOnlineUserCount()); + + + */ + } + + + /** + * 处理订单相关消息(如查询订单、取消订单) + */ + private void handleOrderMsg(WebSocketMessage wsMessage) { + /* + String senderId = wsMessage.getSenderId(); + String data = wsMessage.getData(); + + // 1. 解析订单业务数据 + OrderMsgDTO orderMsg = JSON.parseObject(data, OrderMsgDTO.class); + String orderId = orderMsg.getOrderId(); + String action = orderMsg.getAction(); // 如QUERY/CANCEL + + // 2. 执行业务逻辑 + String result; + if ("QUERY".equals(action)) { + result = orderService.queryOrderById(orderId, senderId); + } else if ("CANCEL".equals(action)) { + result = orderService.cancelOrderById(orderId, senderId); + } else { + result = "不支持的订单操作:" + action; + } + + // 3. 将业务结果推送给发送者 + ChatWebSocketServer.sendToUser(senderId, "【订单处理结果】" + result); + + */ + } + + /** + * 处理心跳包(仅返回确认,无需复杂逻辑) + */ + private void handleHeartbeatMsg(WebSocketMessage wsMessage) { + String senderId = wsMessage.getSenderId(); + // 回复心跳确认 + ChatWebSocketServer.sendToUser(senderId, JSON.toJSONString(new WebSocketMessage("HEARTBEAT_ACK", "pong", senderId, null, System.currentTimeMillis()))); + } +} \ No newline at end of file diff --git a/lis-interfaceCommon/src/main/java/com/czlis/interfaceCommon/websocket/WebSocketUser.java b/lis-interfaceCommon/src/main/java/com/czlis/interfaceCommon/websocket/WebSocketUser.java new file mode 100644 index 0000000..a8284e4 --- /dev/null +++ b/lis-interfaceCommon/src/main/java/com/czlis/interfaceCommon/websocket/WebSocketUser.java @@ -0,0 +1,78 @@ +package com.czlis.interfaceCommon.websocket; + +import javax.websocket.Session; +import java.net.InetSocketAddress; +import java.time.LocalDateTime; +import java.util.Map; +import java.util.concurrent.ConcurrentHashMap; + +/** + * 封装WebSocket用户连接信息:聚合Session + 身份属性 + 连接元数据 + */ +public class WebSocketUser { + // 核心字段:用户唯一标识 + private final String userId; + // 身份属性(可根据业务扩展) + private final String username; // 用户名 + private final String role; // 角色(如ADMIN/USER) + private final String loginIp; // 登录IP + // 连接元数据 + private final Session session; // WebSocket连接会话 + private final LocalDateTime connectTime; // 连接建立时间 + // 动态属性(支持运行时更新,线程安全) + private final Map dynamicAttrs; // 自定义动态属性 + + // 构造器:初始化核心属性(必填) + public WebSocketUser(String userId, String username, String role, String loginIp, Session session) { + this.userId = userId; + this.username = username; + this.role = role; + this.loginIp = loginIp; + this.session = session; + this.connectTime = LocalDateTime.now(); + this.dynamicAttrs = new ConcurrentHashMap<>(); // 动态属性用线程安全Map + } + + // ========== 只读属性(getter) ========== + public String getUserId() { + return userId; + } + + public String getUsername() { + return username; + } + + public String getRole() { + return role; + } + + public String getLoginIp() { + return loginIp; + } + + public Session getSession() { + return session; + } + + public LocalDateTime getConnectTime() { + return connectTime; + } + + // ========== 动态属性(支持增删改查,线程安全) ========== + public void setDynamicAttr(String key, Object value) { + dynamicAttrs.put(key, value); + } + + public Object getDynamicAttr(String key) { + return dynamicAttrs.get(key); + } + + public void removeDynamicAttr(String key) { + dynamicAttrs.remove(key); + } + + // 检查连接是否有效 + public boolean isConnected() { + return session != null && session.isOpen(); + } +} \ No newline at end of file diff --git a/lis-interfaceCommon/src/main/java/com/czlis/interfaceCommon/websocket/notsbook b/lis-interfaceCommon/src/main/java/com/czlis/interfaceCommon/websocket/notsbook new file mode 100644 index 0000000..01cc182 --- /dev/null +++ b/lis-interfaceCommon/src/main/java/com/czlis/interfaceCommon/websocket/notsbook @@ -0,0 +1,57 @@ +websocket消息服务器结构说明 + +1.核心框架 +ChatWebSocketServer +主要作用:提供底层的客户端连接管理,消息的发送与接收,以及客户端身份信息的修改 +WebSocketMessage +说明:消息结构标准实体类 +WebSocketUser +说明:客户端列表属性实体类 +WebSocketMsgService +说明:WebSocket消息处理业务接口 +WebSocketMsgServiceImpl +说明:WebSocket消息处理业务实现类 +RetryMsgService +说明:消息重试服务,发送失败重发,最多3次 +RetryMsgTask +说明:定时器服务器,增加定时业务服务的地方 + +2,协议约定 +消息服务单实例,避免重复实例化错误 +消息服务器地址/ws/msgserver/{SID} +SID为客户端唯一ID,约定使用浏览器指纹 +WebSocketMessage消息结构 + // 消息类型(如CHAT/NOTICE/ORDER/HEARTBEAT) + private String msgType; + // 业务数据(JSON字符串,可解析为对应业务DTO) + private String data; + // 发送者ID(可选,也可从WebSocket连接中获取) + private String senderId; + // 接收者ID(可选,单聊/广播区分) + private String receiverId; + // 消息时间戳 + private Long timestamp; + 消息时间戳必传,目前代码直接默认了 + 重点字段msgType,申明发送的是什么消息 + 发送者ID不需要传,自动获取 + 接收者ID在指定传输目标时传入,目前暂时不使用 + data 字段为消息体,都是json格式 + 消息举例 + {"msgType":"CHAT","data":{"instrid":"abc"},"receiverId":"LIS","timestamp":1755068400000} + {"msgType":"CHAT","data":{"instrid":"abc"}} + 消息类型约定 + 3 业务流程 + 启动服务器,自动创建websocket服务实例 + 客户端启动,连接websocket服务器,发送约定身份参数(非必须) + 客户端等待消息中... + + 服务器创建消息服务业务层 + 定时器注册指定服务为定时任务,并声明循环周期(单线程服务,一类业务只创建一个消息服务器端) + 消息服务程序触发消息规则,给指定客户端发送消息 + 客户端收到消息,做相应的处理 + + 4 应用范例 + UserService + 提供基本的消息服务方法, + 常用的操作客户端方法直接操作ChatWebSocketServer完成 + 包括调阅客户端列表,修改客户端附属属性,过滤,查找,给目标客户端发送消息,以及处理收到的消息