init: 导入RuoYi‑Vue‑Plus 6.X完整代码
This commit is contained in:
@@ -0,0 +1,56 @@
|
||||
<?xml version="1.0" encoding="UTF-8"?>
|
||||
<project xmlns="http://maven.apache.org/POM/4.0.0"
|
||||
xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"
|
||||
xsi:schemaLocation="http://maven.apache.org/POM/4.0.0 http://maven.apache.org/xsd/maven-4.0.0.xsd">
|
||||
<parent>
|
||||
<groupId>org.dromara</groupId>
|
||||
<artifactId>ruoyi-common</artifactId>
|
||||
<version>${revision}</version>
|
||||
</parent>
|
||||
<modelVersion>4.0.0</modelVersion>
|
||||
|
||||
<artifactId>ruoyi-common-push</artifactId>
|
||||
|
||||
<description>
|
||||
ruoyi-common-push 消息推送模块
|
||||
</description>
|
||||
|
||||
<dependencies>
|
||||
<!-- 核心模块 -->
|
||||
<dependency>
|
||||
<groupId>org.dromara</groupId>
|
||||
<artifactId>ruoyi-common-core</artifactId>
|
||||
</dependency>
|
||||
<!-- 缓存服务 -->
|
||||
<dependency>
|
||||
<groupId>org.dromara</groupId>
|
||||
<artifactId>ruoyi-common-redis</artifactId>
|
||||
</dependency>
|
||||
<!-- satoken -->
|
||||
<dependency>
|
||||
<groupId>org.dromara</groupId>
|
||||
<artifactId>ruoyi-common-satoken</artifactId>
|
||||
</dependency>
|
||||
<!-- 序列化模块 -->
|
||||
<dependency>
|
||||
<groupId>org.dromara</groupId>
|
||||
<artifactId>ruoyi-common-json</artifactId>
|
||||
</dependency>
|
||||
<!-- api模块 -->
|
||||
<dependency>
|
||||
<groupId>org.dromara</groupId>
|
||||
<artifactId>ruoyi-api</artifactId>
|
||||
</dependency>
|
||||
<!-- WebSocket 服务 -->
|
||||
<dependency>
|
||||
<groupId>org.springframework.boot</groupId>
|
||||
<artifactId>spring-boot-starter-websocket</artifactId>
|
||||
<exclusions>
|
||||
<exclusion>
|
||||
<groupId>org.springframework.boot</groupId>
|
||||
<artifactId>spring-boot-starter-tomcat</artifactId>
|
||||
</exclusion>
|
||||
</exclusions>
|
||||
</dependency>
|
||||
</dependencies>
|
||||
</project>
|
||||
+24
@@ -0,0 +1,24 @@
|
||||
package org.dromara.common.push.annotation;
|
||||
|
||||
import org.dromara.common.push.condition.MessageTransportCondition;
|
||||
import org.springframework.context.annotation.Conditional;
|
||||
|
||||
import java.lang.annotation.*;
|
||||
|
||||
/**
|
||||
* 按消息推送传输方式启用组件。
|
||||
*
|
||||
* @author Lion Li
|
||||
*/
|
||||
@Documented
|
||||
@Target({ElementType.TYPE, ElementType.METHOD})
|
||||
@Retention(RetentionPolicy.RUNTIME)
|
||||
@Conditional(MessageTransportCondition.class)
|
||||
public @interface ConditionalOnMessageTransport {
|
||||
|
||||
/**
|
||||
* 传输方式:sse / websocket。
|
||||
*/
|
||||
String value();
|
||||
|
||||
}
|
||||
+39
@@ -0,0 +1,39 @@
|
||||
package org.dromara.common.push.condition;
|
||||
|
||||
import org.dromara.common.push.annotation.ConditionalOnMessageTransport;
|
||||
import org.dromara.common.push.enums.MessageTransportEnum;
|
||||
import org.jspecify.annotations.NonNull;
|
||||
import org.springframework.context.annotation.Condition;
|
||||
import org.springframework.context.annotation.ConditionContext;
|
||||
import org.springframework.core.type.AnnotatedTypeMetadata;
|
||||
|
||||
import java.util.Map;
|
||||
|
||||
/**
|
||||
* 消息推送传输方式条件判断。
|
||||
*
|
||||
* @author Lion Li
|
||||
*/
|
||||
public class MessageTransportCondition implements Condition {
|
||||
|
||||
/**
|
||||
* 判断当前消息推送配置是否匹配注解声明的传输方式。
|
||||
*
|
||||
* @param context 条件上下文
|
||||
* @param metadata 注解元数据
|
||||
* @return 是否匹配
|
||||
*/
|
||||
@Override
|
||||
public boolean matches(@NonNull ConditionContext context, AnnotatedTypeMetadata metadata) {
|
||||
Map<String, Object> attributes = metadata.getAnnotationAttributes(ConditionalOnMessageTransport.class.getName());
|
||||
if (attributes == null) {
|
||||
return true;
|
||||
}
|
||||
|
||||
Boolean enabled = context.getEnvironment().getProperty("message.enabled", Boolean.class, true);
|
||||
String transport = context.getEnvironment().getProperty("message.transport", MessageTransportEnum.SSE.getCode());
|
||||
String expected = (String) attributes.get("value");
|
||||
return enabled && expected.equalsIgnoreCase(transport);
|
||||
}
|
||||
|
||||
}
|
||||
+17
@@ -0,0 +1,17 @@
|
||||
package org.dromara.common.push.config;
|
||||
|
||||
import org.dromara.common.push.properties.MessageProperties;
|
||||
import org.springframework.boot.autoconfigure.AutoConfiguration;
|
||||
import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty;
|
||||
import org.springframework.boot.context.properties.EnableConfigurationProperties;
|
||||
|
||||
/**
|
||||
* 统一消息推送公共自动装配。
|
||||
*
|
||||
* @author Lion Li
|
||||
*/
|
||||
@AutoConfiguration
|
||||
@ConditionalOnProperty(prefix = "message", name = "enabled", havingValue = "true", matchIfMissing = true)
|
||||
@EnableConfigurationProperties(MessageProperties.class)
|
||||
public class MessageAutoConfiguration {
|
||||
}
|
||||
+57
@@ -0,0 +1,57 @@
|
||||
package org.dromara.common.push.config;
|
||||
|
||||
import org.dromara.common.push.annotation.ConditionalOnMessageTransport;
|
||||
import org.dromara.common.push.controller.SseController;
|
||||
import org.dromara.common.push.core.SseEmitterSessionManager;
|
||||
import org.dromara.common.push.listener.MessageTopicListener;
|
||||
import org.dromara.common.push.properties.MessageProperties;
|
||||
import org.springframework.boot.autoconfigure.AutoConfiguration;
|
||||
import org.springframework.context.annotation.Bean;
|
||||
|
||||
import java.util.concurrent.ScheduledExecutorService;
|
||||
|
||||
/**
|
||||
* SSE 消息推送自动装配。
|
||||
*
|
||||
* @author Lion Li
|
||||
*/
|
||||
@AutoConfiguration(after = MessageAutoConfiguration.class)
|
||||
@ConditionalOnMessageTransport("sse")
|
||||
public class MessageSseConfiguration {
|
||||
|
||||
/**
|
||||
* 注册 SSE 会话管理器
|
||||
* 负责管理用户 SSE 连接、消息发送、会话清理
|
||||
*
|
||||
* @return SseEmitterSessionManager 实例
|
||||
*/
|
||||
@Bean
|
||||
public SseEmitterSessionManager sseEmitterManager(ScheduledExecutorService scheduledExecutorService,
|
||||
MessageProperties messageProperties) {
|
||||
return new SseEmitterSessionManager(scheduledExecutorService, messageProperties);
|
||||
}
|
||||
|
||||
/**
|
||||
* 注册消息主题监听器
|
||||
* 监听 Redis 全局消息,用于集群环境下的消息分发
|
||||
*
|
||||
* @param manager SSE 会话管理器
|
||||
* @return MessageTopicListener 实例
|
||||
*/
|
||||
@Bean
|
||||
public MessageTopicListener messageTopicListener(SseEmitterSessionManager manager) {
|
||||
return new MessageTopicListener(manager);
|
||||
}
|
||||
|
||||
/**
|
||||
* 注册 SSE 控制器
|
||||
* 提供前端建立 SSE 连接的接口
|
||||
*
|
||||
* @param manager SSE 会话管理器
|
||||
* @return SseController 实例
|
||||
*/
|
||||
@Bean
|
||||
public SseController sseController(SseEmitterSessionManager manager) {
|
||||
return new SseController(manager);
|
||||
}
|
||||
}
|
||||
+79
@@ -0,0 +1,79 @@
|
||||
package org.dromara.common.push.config;
|
||||
|
||||
import org.dromara.common.push.annotation.ConditionalOnMessageTransport;
|
||||
import org.dromara.common.push.core.WebSocketSessionManager;
|
||||
import org.dromara.common.push.handler.PlusWebSocketHandler;
|
||||
import org.dromara.common.push.interceptor.PlusWebSocketInterceptor;
|
||||
import org.dromara.common.push.listener.MessageTopicListener;
|
||||
import org.dromara.common.push.properties.MessageProperties;
|
||||
import org.springframework.boot.autoconfigure.AutoConfiguration;
|
||||
import org.springframework.context.annotation.Bean;
|
||||
import org.springframework.web.socket.WebSocketHandler;
|
||||
import org.springframework.web.socket.config.annotation.EnableWebSocket;
|
||||
import org.springframework.web.socket.config.annotation.WebSocketConfigurer;
|
||||
import org.springframework.web.socket.server.HandshakeInterceptor;
|
||||
|
||||
import java.util.concurrent.ScheduledExecutorService;
|
||||
|
||||
/**
|
||||
* WebSocket 消息推送自动装配。
|
||||
*
|
||||
* @author Lion Li
|
||||
*/
|
||||
@EnableWebSocket
|
||||
@AutoConfiguration(after = MessageAutoConfiguration.class)
|
||||
@ConditionalOnMessageTransport("websocket")
|
||||
public class MessageWebSocketConfiguration {
|
||||
|
||||
/**
|
||||
* WebSocket 配置注册
|
||||
* 配置连接路径、拦截器、跨域
|
||||
*/
|
||||
@Bean
|
||||
public WebSocketConfigurer webSocketConfigurer(HandshakeInterceptor handshakeInterceptor,
|
||||
WebSocketHandler webSocketHandler,
|
||||
MessageProperties messageProperties) {
|
||||
return registry -> registry
|
||||
.addHandler(webSocketHandler, messageProperties.getPath())
|
||||
.addInterceptors(handshakeInterceptor)
|
||||
.setAllowedOrigins(messageProperties.getAllowedOrigins());
|
||||
}
|
||||
|
||||
/**
|
||||
* WebSocket 会话管理器
|
||||
* 负责连接管理、消息发送、定时清理失效会话
|
||||
*/
|
||||
@Bean
|
||||
public WebSocketSessionManager webSocketSessionManager(ScheduledExecutorService scheduledExecutorService,
|
||||
MessageProperties messageProperties) {
|
||||
return new WebSocketSessionManager(scheduledExecutorService, messageProperties);
|
||||
}
|
||||
|
||||
/**
|
||||
* WebSocket 握手拦截器
|
||||
* 建立连接前做登录校验、客户端ID校验
|
||||
*/
|
||||
@Bean
|
||||
public HandshakeInterceptor handshakeInterceptor() {
|
||||
return new PlusWebSocketInterceptor();
|
||||
}
|
||||
|
||||
/**
|
||||
* WebSocket 消息处理器
|
||||
* 处理连接、消息、心跳、断开、异常等事件
|
||||
*/
|
||||
@Bean
|
||||
public WebSocketHandler webSocketHandler(WebSocketSessionManager webSocketSessionManager,
|
||||
MessageProperties messageProperties) {
|
||||
return new PlusWebSocketHandler(webSocketSessionManager, messageProperties);
|
||||
}
|
||||
|
||||
/**
|
||||
* 消息主题监听器
|
||||
* 订阅 Redis 消息,实现集群环境下的消息分发
|
||||
*/
|
||||
@Bean
|
||||
public MessageTopicListener messageTopicListener(WebSocketSessionManager webSocketSessionManager) {
|
||||
return new MessageTopicListener(webSocketSessionManager);
|
||||
}
|
||||
}
|
||||
+39
@@ -0,0 +1,39 @@
|
||||
package org.dromara.common.push.constant;
|
||||
|
||||
/**
|
||||
* 模块通用消息常量定义。
|
||||
*
|
||||
* @author Lion Li
|
||||
*/
|
||||
public interface MessageConstants {
|
||||
|
||||
/**
|
||||
* 登录用户信息
|
||||
*/
|
||||
String LOGIN_USER_KEY = "loginUser";
|
||||
|
||||
/**
|
||||
* 登录令牌
|
||||
*/
|
||||
String LOGIN_TOKEN_KEY = "token";
|
||||
|
||||
/**
|
||||
* 全局消息订阅主题
|
||||
*/
|
||||
String MESSAGE_TOPIC = "global:message";
|
||||
|
||||
/**
|
||||
* 心跳请求标识
|
||||
*/
|
||||
String PING = "ping";
|
||||
|
||||
/**
|
||||
* 心跳响应标识
|
||||
*/
|
||||
String PONG = "pong";
|
||||
|
||||
/**
|
||||
* 同一 token 的新连接替换旧连接时发送给旧连接的控制消息。
|
||||
*/
|
||||
String KICKED = "kicked";
|
||||
}
|
||||
+105
@@ -0,0 +1,105 @@
|
||||
package org.dromara.common.push.controller;
|
||||
|
||||
import cn.dev33.satoken.annotation.SaIgnore;
|
||||
import cn.dev33.satoken.stp.StpUtil;
|
||||
import jakarta.servlet.http.HttpServletResponse;
|
||||
import lombok.RequiredArgsConstructor;
|
||||
import org.dromara.common.core.domain.R;
|
||||
import org.dromara.common.push.annotation.ConditionalOnMessageTransport;
|
||||
import org.dromara.common.push.core.SseEmitterSessionManager;
|
||||
import org.dromara.common.satoken.utils.LoginHelper;
|
||||
import org.springframework.beans.factory.DisposableBean;
|
||||
import org.springframework.http.MediaType;
|
||||
import org.springframework.web.bind.annotation.GetMapping;
|
||||
import org.springframework.web.bind.annotation.RestController;
|
||||
import org.springframework.web.servlet.mvc.method.annotation.SseEmitter;
|
||||
|
||||
/**
|
||||
* SSE 控制器
|
||||
*
|
||||
* @author Lion Li
|
||||
*/
|
||||
@RestController
|
||||
@ConditionalOnMessageTransport("sse")
|
||||
@RequiredArgsConstructor
|
||||
public class SseController implements DisposableBean {
|
||||
|
||||
private final SseEmitterSessionManager sessionManager;
|
||||
|
||||
/**
|
||||
* 建立当前登录用户的 SSE 连接。
|
||||
*
|
||||
* @return SSE 发射器
|
||||
*/
|
||||
@GetMapping(value = "${message.path:/resource/message}", produces = MediaType.TEXT_EVENT_STREAM_VALUE)
|
||||
public SseEmitter connect(HttpServletResponse response) {
|
||||
prepareSseResponse(response);
|
||||
String tokenValue = StpUtil.getTokenValue();
|
||||
Long userId = LoginHelper.getUserId();
|
||||
return sessionManager.connect(userId, tokenValue);
|
||||
}
|
||||
|
||||
/**
|
||||
* 关闭当前登录用户的 SSE 连接。
|
||||
*
|
||||
* @return 操作结果
|
||||
*/
|
||||
@SaIgnore
|
||||
@GetMapping(value = "${message.path:/resource/message}/close")
|
||||
public R<Void> close() {
|
||||
String tokenValue = StpUtil.getTokenValue();
|
||||
Long userId = LoginHelper.getUserId();
|
||||
sessionManager.disconnect(userId, tokenValue);
|
||||
return R.ok();
|
||||
}
|
||||
|
||||
/**
|
||||
* 设置 SSE 响应头,覆盖统一鉴权成功路径中的默认 JSON 响应类型。
|
||||
*
|
||||
* @param response 当前响应
|
||||
*/
|
||||
private void prepareSseResponse(HttpServletResponse response) {
|
||||
response.setContentType(MediaType.TEXT_EVENT_STREAM_VALUE);
|
||||
response.setCharacterEncoding("UTF-8");
|
||||
response.setHeader("Cache-Control", "no-cache");
|
||||
response.setHeader("X-Accel-Buffering", "no");
|
||||
}
|
||||
|
||||
// 以下为demo仅供参考 禁止使用 请在业务逻辑中使用工具发送而不是用接口发送
|
||||
// /**
|
||||
// * 向特定用户发送消息
|
||||
// *
|
||||
// * @param userId 目标用户的 ID
|
||||
// * @param msg 要发送的消息内容
|
||||
// */
|
||||
// @GetMapping(value = "${message.path:/resource/message}/send")
|
||||
// public R<Void> send(Long userId, String msg) {
|
||||
// PushDTO dto = new PushDTO();
|
||||
// dto.setUserIds(List.of(userId));
|
||||
// dto.setPayload(PushPayloadDTO.of("message", "backend", msg, null));
|
||||
// sessionManager.publishMessage(dto);
|
||||
// return R.ok();
|
||||
// }
|
||||
//
|
||||
// /**
|
||||
// * 向所有用户发送消息
|
||||
// *
|
||||
// * @param msg 要发送的消息内容
|
||||
// */
|
||||
// @GetMapping(value = "${message.path:/resource/message}/sendAll")
|
||||
// public R<Void> send(String msg) {
|
||||
// sessionManager.publishAll(msg);
|
||||
// return R.ok();
|
||||
// }
|
||||
|
||||
/**
|
||||
* 容器销毁时释放资源占位实现。
|
||||
*
|
||||
* @throws Exception 销毁异常
|
||||
*/
|
||||
@Override
|
||||
public void destroy() throws Exception {
|
||||
// 销毁时不需要做什么 此方法避免无用操作报错
|
||||
}
|
||||
|
||||
}
|
||||
+51
@@ -0,0 +1,51 @@
|
||||
package org.dromara.common.push.core;
|
||||
|
||||
import org.dromara.common.push.dto.PushDTO;
|
||||
import org.dromara.system.api.domain.PushPayloadDTO;
|
||||
|
||||
import java.util.function.Consumer;
|
||||
|
||||
/**
|
||||
* 统一推送会话管理器。
|
||||
*
|
||||
* @author Lion Li
|
||||
*/
|
||||
public interface PushSessionManager {
|
||||
|
||||
/**
|
||||
* 订阅消息通道
|
||||
* 注册消息消费者,用于监听并处理消息推送事件
|
||||
*
|
||||
* @param consumer 消息消费逻辑
|
||||
*/
|
||||
void subscribeMessage(Consumer<PushDTO> consumer);
|
||||
|
||||
/**
|
||||
* 发送消息给指定用户
|
||||
*
|
||||
* @param userId 目标用户ID
|
||||
* @param payload 消息体
|
||||
*/
|
||||
void sendMessage(Long userId, PushPayloadDTO payload);
|
||||
|
||||
/**
|
||||
* 全局广播消息(所有在线用户)
|
||||
*
|
||||
* @param payload 消息体
|
||||
*/
|
||||
void sendMessage(PushPayloadDTO payload);
|
||||
|
||||
/**
|
||||
* 批量发布消息给指定用户列表
|
||||
*
|
||||
* @param pushDTO 推送参数封装对象
|
||||
*/
|
||||
void publishMessage(PushDTO pushDTO);
|
||||
|
||||
/**
|
||||
* 全局广播消息(所有用户)
|
||||
*
|
||||
* @param payload 消息体
|
||||
*/
|
||||
void publishAll(PushPayloadDTO payload);
|
||||
}
|
||||
+309
@@ -0,0 +1,309 @@
|
||||
package org.dromara.common.push.core;
|
||||
|
||||
import cn.hutool.core.collection.CollUtil;
|
||||
import cn.hutool.core.map.MapUtil;
|
||||
import lombok.extern.slf4j.Slf4j;
|
||||
import org.dromara.common.core.utils.ThreadUtils;
|
||||
import org.dromara.common.json.utils.JsonUtils;
|
||||
import org.dromara.common.push.constant.MessageConstants;
|
||||
import org.dromara.common.push.dto.PushDTO;
|
||||
import org.dromara.common.push.properties.MessageProperties;
|
||||
import org.dromara.common.redis.utils.RedisUtils;
|
||||
import org.dromara.system.api.domain.PushPayloadDTO;
|
||||
import org.springframework.web.servlet.mvc.method.annotation.SseEmitter;
|
||||
|
||||
import java.io.IOException;
|
||||
import java.util.ArrayList;
|
||||
import java.util.List;
|
||||
import java.util.Map;
|
||||
import java.util.concurrent.ConcurrentHashMap;
|
||||
import java.util.concurrent.ScheduledExecutorService;
|
||||
import java.util.concurrent.TimeUnit;
|
||||
import java.util.function.Consumer;
|
||||
|
||||
/**
|
||||
* 管理 Server-Sent Events (SSE) 连接
|
||||
*
|
||||
* @author Lion Li
|
||||
*/
|
||||
@Slf4j
|
||||
public class SseEmitterSessionManager implements PushSessionManager {
|
||||
|
||||
private final static Map<Long, Map<String, SseEmitter>> USER_TOKEN_EMITTERS = new ConcurrentHashMap<>();
|
||||
|
||||
private final MessageProperties messageProperties;
|
||||
|
||||
/**
|
||||
* 构造 SSE 会话管理器并启动心跳检测。
|
||||
*
|
||||
* @param scheduledExecutorService 定时任务线程池
|
||||
* @param messageProperties 消息推送配置
|
||||
*/
|
||||
public SseEmitterSessionManager(ScheduledExecutorService scheduledExecutorService, MessageProperties messageProperties) {
|
||||
this.messageProperties = messageProperties;
|
||||
// 定时执行 SSE 心跳检测
|
||||
scheduledExecutorService.scheduleWithFixedDelay(
|
||||
this::sseMonitor,
|
||||
messageProperties.getHeartbeatInterval(),
|
||||
messageProperties.getHeartbeatInterval(),
|
||||
TimeUnit.SECONDS
|
||||
);
|
||||
}
|
||||
|
||||
/**
|
||||
* 建立与指定用户的 SSE 连接
|
||||
*
|
||||
* @param userId 用户的唯一标识符,用于区分不同用户的连接
|
||||
* @param token 用户的唯一令牌,用于识别具体的连接
|
||||
* @return 返回一个 SseEmitter 实例,客户端可以通过该实例接收 SSE 事件
|
||||
*/
|
||||
public SseEmitter connect(Long userId, String token) {
|
||||
// 从 USER_TOKEN_EMITTERS 中获取或创建当前用户的 SseEmitter 映射表(ConcurrentHashMap)
|
||||
// 每个用户可以有多个 SSE 连接,通过 token 进行区分
|
||||
Map<String, SseEmitter> emitters = USER_TOKEN_EMITTERS.computeIfAbsent(userId, k -> new ConcurrentHashMap<>());
|
||||
|
||||
// 关闭已存在的SseEmitter,防止超过最大连接数
|
||||
SseEmitter oldEmitter = emitters.remove(token);
|
||||
if (oldEmitter != null) {
|
||||
sendKickedMessage(oldEmitter);
|
||||
oldEmitter.complete();
|
||||
}
|
||||
|
||||
// 创建一个新的 SseEmitter 实例,避免连接之后直接关闭浏览器导致连接停滞
|
||||
SseEmitter emitter = new SseEmitter(messageProperties.getSseTimeout());
|
||||
|
||||
emitters.put(token, emitter);
|
||||
|
||||
// 当 emitter 完成、超时或发生错误时,从映射表中移除对应的 token
|
||||
emitter.onCompletion(() -> {
|
||||
SseEmitter remove = emitters.remove(token);
|
||||
if (remove != null) {
|
||||
remove.complete();
|
||||
}
|
||||
});
|
||||
emitter.onTimeout(() -> {
|
||||
SseEmitter remove = emitters.remove(token);
|
||||
if (remove != null) {
|
||||
remove.complete();
|
||||
}
|
||||
});
|
||||
emitter.onError((e) -> {
|
||||
SseEmitter remove = emitters.remove(token);
|
||||
if (remove != null) {
|
||||
remove.complete();
|
||||
}
|
||||
});
|
||||
|
||||
try {
|
||||
// 向客户端发送一条连接成功的事件
|
||||
emitter.send(SseEmitter.event().comment("connected"));
|
||||
} catch (IOException e) {
|
||||
// 如果发送消息失败,则从映射表中移除 emitter
|
||||
emitters.remove(token);
|
||||
}
|
||||
return emitter;
|
||||
}
|
||||
|
||||
/**
|
||||
* 通知旧连接已被同 token 新连接替换。
|
||||
*
|
||||
* @param emitter 旧 SSE 连接
|
||||
*/
|
||||
private void sendKickedMessage(SseEmitter emitter) {
|
||||
try {
|
||||
emitter.send(SseEmitter.event()
|
||||
.name("message")
|
||||
.data(MessageConstants.KICKED));
|
||||
} catch (Exception ignore) {
|
||||
// 旧连接可能已断开,忽略通知失败
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* 断开指定用户的 SSE 连接
|
||||
*
|
||||
* @param userId 用户的唯一标识符,用于区分不同用户的连接
|
||||
* @param token 用户的唯一令牌,用于识别具体的连接
|
||||
*/
|
||||
public void disconnect(Long userId, String token) {
|
||||
if (userId == null || token == null) {
|
||||
return;
|
||||
}
|
||||
Map<String, SseEmitter> emitters = USER_TOKEN_EMITTERS.get(userId);
|
||||
if (MapUtil.isNotEmpty(emitters)) {
|
||||
try {
|
||||
SseEmitter sseEmitter = emitters.get(token);
|
||||
sseEmitter.send(SseEmitter.event().comment("disconnected"));
|
||||
sseEmitter.complete();
|
||||
} catch (Exception ignore) {
|
||||
}
|
||||
emitters.remove(token);
|
||||
} else {
|
||||
USER_TOKEN_EMITTERS.remove(userId);
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* 执行 SSE 心跳检测并清理失效连接。
|
||||
*/
|
||||
public void sseMonitor() {
|
||||
final SseEmitter.SseEventBuilder heartbeat = SseEmitter.event().comment("heartbeat");
|
||||
// 记录需要移除的用户ID
|
||||
List<Long> toRemoveUsers = new ArrayList<>();
|
||||
|
||||
USER_TOKEN_EMITTERS.forEach((userId, emitterMap) -> {
|
||||
if (CollUtil.isEmpty(emitterMap)) {
|
||||
toRemoveUsers.add(userId);
|
||||
return;
|
||||
}
|
||||
|
||||
emitterMap.entrySet().removeIf(entry -> {
|
||||
try {
|
||||
entry.getValue().send(heartbeat);
|
||||
return false;
|
||||
} catch (Exception ex) {
|
||||
try {
|
||||
entry.getValue().complete();
|
||||
} catch (Exception ignore) {
|
||||
// 忽略重复关闭异常
|
||||
}
|
||||
return true; // 发送失败 → 移除该连接
|
||||
}
|
||||
});
|
||||
|
||||
// 移除空连接用户
|
||||
if (emitterMap.isEmpty()) {
|
||||
toRemoveUsers.add(userId);
|
||||
}
|
||||
});
|
||||
|
||||
// 循环结束后统一清理空用户,避免并发修改异常
|
||||
toRemoveUsers.forEach(USER_TOKEN_EMITTERS::remove);
|
||||
}
|
||||
|
||||
/**
|
||||
* 订阅 SSE 广播主题消息。
|
||||
*
|
||||
* @param consumer 处理SSE消息的消费者函数
|
||||
*/
|
||||
@Override
|
||||
public void subscribeMessage(Consumer<PushDTO> consumer) {
|
||||
RedisUtils.subscribe(MessageConstants.MESSAGE_TOPIC, PushDTO.class, consumer);
|
||||
}
|
||||
|
||||
/**
|
||||
* 向指定用户的全部本地 SSE 会话发送消息。
|
||||
*
|
||||
* @param userId 要发送消息的用户id
|
||||
* @param message 要发送的消息内容
|
||||
*/
|
||||
public void sendMessage(Long userId, String message) {
|
||||
Map<String, SseEmitter> emitters = USER_TOKEN_EMITTERS.get(userId);
|
||||
if (MapUtil.isNotEmpty(emitters)) {
|
||||
for (Map.Entry<String, SseEmitter> entry : emitters.entrySet()) {
|
||||
try {
|
||||
entry.getValue().send(SseEmitter.event()
|
||||
.name("message")
|
||||
.data(message));
|
||||
} catch (Exception e) {
|
||||
SseEmitter remove = emitters.remove(entry.getKey());
|
||||
if (remove != null) {
|
||||
remove.complete();
|
||||
}
|
||||
}
|
||||
}
|
||||
} else {
|
||||
USER_TOKEN_EMITTERS.remove(userId);
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* 向指定用户的全部本地 SSE 会话发送统一 JSON 消息。
|
||||
*
|
||||
* @param userId 要发送消息的用户id
|
||||
* @param payload 要发送的消息体
|
||||
*/
|
||||
@Override
|
||||
public void sendMessage(Long userId, PushPayloadDTO payload) {
|
||||
if (payload == null) {
|
||||
return;
|
||||
}
|
||||
sendMessage(userId, JsonUtils.toJsonString(payload));
|
||||
}
|
||||
|
||||
/**
|
||||
* 向指定用户的全部本地 SSE 会话发送统一 JSON 消息。
|
||||
*
|
||||
* @param userId 要发送消息的用户id
|
||||
* @param pushDTO 要发送的消息内容
|
||||
*/
|
||||
public void sendMessage(Long userId, PushDTO pushDTO) {
|
||||
if (pushDTO == null) {
|
||||
return;
|
||||
}
|
||||
sendMessage(userId, pushDTO.getPayload());
|
||||
}
|
||||
|
||||
/**
|
||||
* 向当前节点所有 SSE 会话发送消息。
|
||||
*
|
||||
* @param message 要发送的消息内容
|
||||
*/
|
||||
public void sendMessage(String message) {
|
||||
List<Long> userIds = new ArrayList<>(USER_TOKEN_EMITTERS.keySet());
|
||||
Runnable[] sendTasks = userIds.stream()
|
||||
.map(userId -> (Runnable) () -> sendMessage(userId, message))
|
||||
.toArray(Runnable[]::new);
|
||||
ThreadUtils.virtualInvokeAll(sendTasks);
|
||||
}
|
||||
|
||||
/**
|
||||
* 向当前节点所有 SSE 会话发送统一 JSON 消息。
|
||||
*
|
||||
* @param payload 要发送的消息体
|
||||
*/
|
||||
@Override
|
||||
public void sendMessage(PushPayloadDTO payload) {
|
||||
if (payload == null) {
|
||||
return;
|
||||
}
|
||||
sendMessage(JsonUtils.toJsonString(payload));
|
||||
}
|
||||
|
||||
/**
|
||||
* 发布 SSE 订阅消息。
|
||||
*
|
||||
* @param pushDTO 要发布的SSE消息对象
|
||||
*/
|
||||
@Override
|
||||
public void publishMessage(PushDTO pushDTO) {
|
||||
if (pushDTO == null || pushDTO.getPayload() == null) {
|
||||
return;
|
||||
}
|
||||
RedisUtils.publish(MessageConstants.MESSAGE_TOPIC, pushDTO, consumer -> log.info(
|
||||
"发送主题订阅消息topic:{} userIds:{} message:{}",
|
||||
MessageConstants.MESSAGE_TOPIC,
|
||||
pushDTO.getUserIds(),
|
||||
pushDTO.getPayload() == null ? null : pushDTO.getPayload().getMessage()
|
||||
));
|
||||
}
|
||||
|
||||
/**
|
||||
* 发布 SSE 广播消息。
|
||||
*
|
||||
* @param message 要发布的消息内容
|
||||
*/
|
||||
public void publishAll(String message) {
|
||||
publishAll(PushPayloadDTO.of("message", "backend", message, null));
|
||||
}
|
||||
|
||||
/**
|
||||
* 发布 SSE 广播 JSON 消息。
|
||||
*
|
||||
* @param payload 要发布的消息体
|
||||
*/
|
||||
@Override
|
||||
public void publishAll(PushPayloadDTO payload) {
|
||||
publishMessage(PushDTO.broadcast(payload));
|
||||
}
|
||||
}
|
||||
+292
@@ -0,0 +1,292 @@
|
||||
package org.dromara.common.push.core;
|
||||
|
||||
import cn.hutool.core.collection.CollUtil;
|
||||
import cn.hutool.core.map.MapUtil;
|
||||
import lombok.extern.slf4j.Slf4j;
|
||||
import org.dromara.common.core.utils.ThreadUtils;
|
||||
import org.dromara.common.json.utils.JsonUtils;
|
||||
import org.dromara.common.push.constant.MessageConstants;
|
||||
import org.dromara.common.push.dto.PushDTO;
|
||||
import org.dromara.common.push.properties.MessageProperties;
|
||||
import org.dromara.common.redis.utils.RedisUtils;
|
||||
import org.dromara.system.api.domain.PushPayloadDTO;
|
||||
import org.springframework.web.socket.*;
|
||||
|
||||
import java.io.IOException;
|
||||
import java.util.ArrayList;
|
||||
import java.util.List;
|
||||
import java.util.Map;
|
||||
import java.util.concurrent.ConcurrentHashMap;
|
||||
import java.util.concurrent.ScheduledExecutorService;
|
||||
import java.util.concurrent.TimeUnit;
|
||||
import java.util.function.Consumer;
|
||||
|
||||
import static org.dromara.common.push.constant.MessageConstants.MESSAGE_TOPIC;
|
||||
|
||||
/**
|
||||
* WebSocket 会话管理器。
|
||||
*
|
||||
* @author Lion Li
|
||||
*/
|
||||
@Slf4j
|
||||
public class WebSocketSessionManager implements PushSessionManager {
|
||||
|
||||
/**
|
||||
* 用户会话存储集合
|
||||
* 结构:userId -> (token -> WebSocketSession)
|
||||
* 支持同一用户多终端、多设备同时在线
|
||||
*/
|
||||
private static final Map<Long, Map<String, WebSocketSession>> USER_TOKEN_SESSIONS = new ConcurrentHashMap<>();
|
||||
|
||||
/**
|
||||
* 构造函数
|
||||
* 初始化定时任务:每60秒执行一次会话监控,自动清理无效连接
|
||||
*/
|
||||
public WebSocketSessionManager(ScheduledExecutorService scheduledExecutorService, MessageProperties messageProperties) {
|
||||
scheduledExecutorService.scheduleWithFixedDelay(
|
||||
this::sessionMonitor,
|
||||
messageProperties.getHeartbeatInterval(),
|
||||
messageProperties.getHeartbeatInterval(),
|
||||
TimeUnit.SECONDS
|
||||
);
|
||||
}
|
||||
|
||||
/**
|
||||
* 用户建立WebSocket连接
|
||||
*
|
||||
* @param userId 用户ID
|
||||
* @param token 客户端唯一标识(区分不同设备/终端)
|
||||
* @param session WebSocket会话对象
|
||||
*/
|
||||
public void connect(Long userId, String token, WebSocketSession session) {
|
||||
Map<String, WebSocketSession> sessions = USER_TOKEN_SESSIONS.computeIfAbsent(userId, key -> new ConcurrentHashMap<>());
|
||||
// 移除并关闭旧的同token会话,避免重复连接
|
||||
WebSocketSession oldSession = sessions.remove(token);
|
||||
sendKickedMessage(oldSession);
|
||||
closeSession(oldSession, CloseStatus.NORMAL);
|
||||
// 存储新会话
|
||||
sessions.put(token, session);
|
||||
}
|
||||
|
||||
/**
|
||||
* 用户断开WebSocket连接
|
||||
*
|
||||
* @param userId 用户ID
|
||||
* @param token 客户端唯一标识
|
||||
*/
|
||||
public void disconnect(Long userId, String token) {
|
||||
disconnect(userId, token, CloseStatus.NORMAL);
|
||||
}
|
||||
|
||||
/**
|
||||
* 用户断开WebSocket连接
|
||||
*
|
||||
* @param userId 用户ID
|
||||
* @param token 客户端唯一标识
|
||||
* @param status 关闭状态码
|
||||
*/
|
||||
public void disconnect(Long userId, String token, CloseStatus status) {
|
||||
if (userId == null || token == null) {
|
||||
return;
|
||||
}
|
||||
Map<String, WebSocketSession> sessions = USER_TOKEN_SESSIONS.get(userId);
|
||||
if (MapUtil.isEmpty(sessions)) {
|
||||
USER_TOKEN_SESSIONS.remove(userId);
|
||||
return;
|
||||
}
|
||||
// 移除指定token会话并关闭
|
||||
closeSession(sessions.remove(token), status);
|
||||
// 该用户无任何会话时,从缓存中移除
|
||||
if (sessions.isEmpty()) {
|
||||
USER_TOKEN_SESSIONS.remove(userId);
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* 会话监控定时任务
|
||||
* 定期清理已关闭、失效的WebSocket会话,防止内存泄漏
|
||||
*/
|
||||
public void sessionMonitor() {
|
||||
List<Long> toRemoveUsers = new ArrayList<>();
|
||||
USER_TOKEN_SESSIONS.forEach((userId, sessionMap) -> {
|
||||
if (CollUtil.isEmpty(sessionMap)) {
|
||||
toRemoveUsers.add(userId);
|
||||
return;
|
||||
}
|
||||
// 移除已关闭的无效会话
|
||||
sessionMap.entrySet().removeIf(entry -> {
|
||||
WebSocketSession session = entry.getValue();
|
||||
if (session == null || !session.isOpen()) {
|
||||
closeSession(session, CloseStatus.NORMAL);
|
||||
return true;
|
||||
}
|
||||
return false;
|
||||
});
|
||||
// 无有效会话,标记用户待删除
|
||||
if (sessionMap.isEmpty()) {
|
||||
toRemoveUsers.add(userId);
|
||||
}
|
||||
});
|
||||
// 批量清理无会话用户
|
||||
toRemoveUsers.forEach(USER_TOKEN_SESSIONS::remove);
|
||||
}
|
||||
|
||||
/**
|
||||
* 通知旧连接已被同 token 新连接替换。
|
||||
*
|
||||
* @param session 旧 WebSocket 会话
|
||||
*/
|
||||
private void sendKickedMessage(WebSocketSession session) {
|
||||
if (session == null || !session.isOpen()) {
|
||||
return;
|
||||
}
|
||||
sendMessage(session, MessageConstants.KICKED);
|
||||
}
|
||||
|
||||
/**
|
||||
* 订阅消息通道
|
||||
* 注册消息消费者,监听Redis消息推送
|
||||
*
|
||||
* @param consumer 消息消费逻辑
|
||||
*/
|
||||
@Override
|
||||
public void subscribeMessage(Consumer<PushDTO> consumer) {
|
||||
RedisUtils.subscribe(MESSAGE_TOPIC, PushDTO.class, consumer);
|
||||
}
|
||||
|
||||
/**
|
||||
* 向指定用户发送消息
|
||||
*
|
||||
* @param userId 目标用户ID
|
||||
* @param payload 消息体
|
||||
*/
|
||||
@Override
|
||||
public void sendMessage(Long userId, PushPayloadDTO payload) {
|
||||
if (payload == null) {
|
||||
return;
|
||||
}
|
||||
Map<String, WebSocketSession> sessions = USER_TOKEN_SESSIONS.get(userId);
|
||||
if (MapUtil.isEmpty(sessions)) {
|
||||
USER_TOKEN_SESSIONS.remove(userId);
|
||||
return;
|
||||
}
|
||||
// 发送消息并自动清理失效会话
|
||||
sessions.entrySet().removeIf(entry -> {
|
||||
WebSocketSession session = entry.getValue();
|
||||
if (session == null || !session.isOpen()) {
|
||||
closeSession(session, CloseStatus.NORMAL);
|
||||
return true;
|
||||
}
|
||||
// 发送失败的会话也会被移除
|
||||
return !sendMessage(session, new TextMessage(JsonUtils.toJsonString(payload)));
|
||||
});
|
||||
// 无有效会话则移除用户
|
||||
if (sessions.isEmpty()) {
|
||||
USER_TOKEN_SESSIONS.remove(userId);
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* 向所有在线用户广播消息
|
||||
*
|
||||
* @param payload 消息体
|
||||
*/
|
||||
@Override
|
||||
public void sendMessage(PushPayloadDTO payload) {
|
||||
if (payload == null) {
|
||||
return;
|
||||
}
|
||||
List<Long> userIds = new ArrayList<>(USER_TOKEN_SESSIONS.keySet());
|
||||
Runnable[] sendTasks = userIds.stream()
|
||||
.map(userId -> (Runnable) () -> sendMessage(userId, payload))
|
||||
.toArray(Runnable[]::new);
|
||||
ThreadUtils.virtualInvokeAll(sendTasks);
|
||||
}
|
||||
|
||||
/**
|
||||
* 发布消息到Redis订阅通道
|
||||
* 支持集群环境下的分布式消息推送
|
||||
*
|
||||
* @param pushDTO 推送消息封装对象
|
||||
*/
|
||||
@Override
|
||||
public void publishMessage(PushDTO pushDTO) {
|
||||
if (pushDTO == null || pushDTO.getPayload() == null) {
|
||||
return;
|
||||
}
|
||||
RedisUtils.publish(MESSAGE_TOPIC, pushDTO, consumer -> log.info(
|
||||
"WebSocket发送主题订阅消息topic:{} userIds:{} message:{}",
|
||||
MESSAGE_TOPIC,
|
||||
pushDTO.getUserIds(),
|
||||
pushDTO.getPayload() == null ? null : pushDTO.getPayload().getMessage()
|
||||
));
|
||||
}
|
||||
|
||||
/**
|
||||
* 全局广播消息(所有用户)
|
||||
*
|
||||
* @param payload 消息体
|
||||
*/
|
||||
@Override
|
||||
public void publishAll(PushPayloadDTO payload) {
|
||||
publishMessage(PushDTO.broadcast(payload));
|
||||
}
|
||||
|
||||
/**
|
||||
* 发送心跳Pong消息
|
||||
* 用于维持WebSocket长连接存活
|
||||
*
|
||||
* @param session WebSocket会话
|
||||
*/
|
||||
public void sendPongMessage(WebSocketSession session) {
|
||||
sendMessage(session, new PongMessage());
|
||||
}
|
||||
|
||||
/**
|
||||
* 发送文本消息
|
||||
*
|
||||
* @param session WebSocket会话
|
||||
* @param message 文本内容
|
||||
*/
|
||||
public void sendMessage(WebSocketSession session, String message) {
|
||||
sendMessage(session, new TextMessage(message));
|
||||
}
|
||||
|
||||
/**
|
||||
* 底层消息发送方法
|
||||
*
|
||||
* @param session 会话对象
|
||||
* @param message WebSocket消息对象
|
||||
* @return 发送是否成功
|
||||
*/
|
||||
private boolean sendMessage(WebSocketSession session, WebSocketMessage<?> message) {
|
||||
if (session == null || !session.isOpen()) {
|
||||
log.warn("[send] session会话已经关闭");
|
||||
return false;
|
||||
}
|
||||
try {
|
||||
session.sendMessage(message);
|
||||
return true;
|
||||
} catch (IOException e) {
|
||||
log.error("[send] session({}) 发送消息({}) 异常", session, message, e);
|
||||
return false;
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* 安全关闭WebSocket会话
|
||||
*
|
||||
* @param session 待关闭的会话
|
||||
* @param status 关闭状态码
|
||||
*/
|
||||
public void closeSession(WebSocketSession session, CloseStatus status) {
|
||||
if (session == null) {
|
||||
return;
|
||||
}
|
||||
try {
|
||||
session.close(status);
|
||||
} catch (Exception ignored) {
|
||||
// 关闭异常忽略,防止影响主流程
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,54 @@
|
||||
package org.dromara.common.push.dto;
|
||||
|
||||
import lombok.Data;
|
||||
import org.dromara.system.api.domain.PushPayloadDTO;
|
||||
|
||||
import java.io.Serial;
|
||||
import java.io.Serializable;
|
||||
import java.util.List;
|
||||
|
||||
/**
|
||||
* 统一推送 DTO。
|
||||
*
|
||||
* @author Lion Li
|
||||
*/
|
||||
@Data
|
||||
public class PushDTO implements Serializable {
|
||||
|
||||
@Serial
|
||||
private static final long serialVersionUID = 1L;
|
||||
|
||||
/**
|
||||
* 目标用户 ID 列表,为空表示广播。
|
||||
*/
|
||||
private List<Long> userIds;
|
||||
|
||||
/**
|
||||
* 推送消息体。
|
||||
*/
|
||||
private PushPayloadDTO payload;
|
||||
|
||||
/**
|
||||
* 构建指定用户推送消息。
|
||||
*
|
||||
* @param userIds 目标用户 ID 列表
|
||||
* @param payload 推送消息体
|
||||
* @return 推送 DTO
|
||||
*/
|
||||
public static PushDTO of(List<Long> userIds, PushPayloadDTO payload) {
|
||||
PushDTO dto = new PushDTO();
|
||||
dto.setUserIds(userIds);
|
||||
dto.setPayload(payload);
|
||||
return dto;
|
||||
}
|
||||
|
||||
/**
|
||||
* 构建广播推送消息。
|
||||
*
|
||||
* @param payload 推送消息体
|
||||
* @return 推送 DTO
|
||||
*/
|
||||
public static PushDTO broadcast(PushPayloadDTO payload) {
|
||||
return of(null, payload);
|
||||
}
|
||||
}
|
||||
+57
@@ -0,0 +1,57 @@
|
||||
package org.dromara.common.push.enums;
|
||||
|
||||
import lombok.AllArgsConstructor;
|
||||
import lombok.Getter;
|
||||
|
||||
import java.util.Arrays;
|
||||
|
||||
/**
|
||||
* 消息推送传输方式。
|
||||
*
|
||||
* @author Lion Li
|
||||
*/
|
||||
@Getter
|
||||
@AllArgsConstructor
|
||||
public enum MessageTransportEnum {
|
||||
|
||||
/**
|
||||
* SSE 传输方式
|
||||
* 服务端推送事件,单向轻量传输
|
||||
*/
|
||||
SSE("sse"),
|
||||
|
||||
/**
|
||||
* WebSocket 传输方式
|
||||
* 全双工长连接,支持双向实时通信
|
||||
*/
|
||||
WEBSOCKET("websocket");
|
||||
|
||||
/**
|
||||
* 传输类型编码
|
||||
*/
|
||||
private final String code;
|
||||
|
||||
/**
|
||||
* 判断传输方式是否匹配
|
||||
*
|
||||
* @param transport 传输方式字符串
|
||||
* @return 是否匹配
|
||||
*/
|
||||
public boolean matches(String transport) {
|
||||
return code.equalsIgnoreCase(transport);
|
||||
}
|
||||
|
||||
/**
|
||||
* 根据传输类型字符串获取枚举
|
||||
* 找不到则默认返回 SSE
|
||||
*
|
||||
* @param transport 传输方式字符串
|
||||
* @return 对应的消息传输枚举
|
||||
*/
|
||||
public static MessageTransportEnum of(String transport) {
|
||||
return Arrays.stream(values())
|
||||
.filter(item -> item.matches(transport))
|
||||
.findFirst()
|
||||
.orElse(SSE);
|
||||
}
|
||||
}
|
||||
+164
@@ -0,0 +1,164 @@
|
||||
package org.dromara.common.push.handler;
|
||||
|
||||
import cn.hutool.core.util.ObjectUtil;
|
||||
import lombok.RequiredArgsConstructor;
|
||||
import lombok.extern.slf4j.Slf4j;
|
||||
import org.dromara.common.core.enums.PushSourceEnum;
|
||||
import org.dromara.common.core.enums.PushTypeEnum;
|
||||
import org.dromara.common.core.utils.StringUtils;
|
||||
import org.dromara.common.push.constant.MessageConstants;
|
||||
import org.dromara.common.push.core.WebSocketSessionManager;
|
||||
import org.dromara.common.push.dto.PushDTO;
|
||||
import org.dromara.common.push.properties.MessageProperties;
|
||||
import org.dromara.system.api.domain.PushPayloadDTO;
|
||||
import org.dromara.system.api.model.LoginUser;
|
||||
import org.springframework.web.socket.*;
|
||||
import org.springframework.web.socket.handler.AbstractWebSocketHandler;
|
||||
import org.springframework.web.socket.handler.ConcurrentWebSocketSessionDecorator;
|
||||
|
||||
import java.io.IOException;
|
||||
import java.util.List;
|
||||
|
||||
/**
|
||||
* WebSocket 请求处理器
|
||||
* 处理WebSocket连接建立、消息接收、异常、断开等全生命周期事件
|
||||
*
|
||||
* @author Lion Li
|
||||
*/
|
||||
@RequiredArgsConstructor
|
||||
@Slf4j
|
||||
public class PlusWebSocketHandler extends AbstractWebSocketHandler {
|
||||
|
||||
/**
|
||||
* WebSocket 会话管理器
|
||||
*/
|
||||
private final WebSocketSessionManager webSocketSessionManager;
|
||||
|
||||
/**
|
||||
* 消息推送配置
|
||||
*/
|
||||
private final MessageProperties messageProperties;
|
||||
|
||||
/**
|
||||
* 建立WebSocket连接后触发
|
||||
* 校验用户登录信息,注册会话
|
||||
*
|
||||
* @param session WebSocket会话
|
||||
*/
|
||||
@Override
|
||||
public void afterConnectionEstablished(WebSocketSession session) throws IOException {
|
||||
// 从会话属性中获取登录用户信息和Token
|
||||
LoginUser loginUser = (LoginUser) session.getAttributes().get(MessageConstants.LOGIN_USER_KEY);
|
||||
String token = (String) session.getAttributes().get(MessageConstants.LOGIN_TOKEN_KEY);
|
||||
|
||||
// 校验用户信息是否为空,无效则直接关闭连接
|
||||
if (ObjectUtil.hasNull(loginUser, token)) {
|
||||
session.close(CloseStatus.BAD_DATA);
|
||||
log.info("[connect] invalid token received. sessionId: {}", session.getId());
|
||||
return;
|
||||
}
|
||||
|
||||
// 并发安全包装会话,并注册到会话管理器
|
||||
webSocketSessionManager.connect(
|
||||
loginUser.getUserId(),
|
||||
token,
|
||||
new ConcurrentWebSocketSessionDecorator(
|
||||
session,
|
||||
messageProperties.getWebSocketSendTimeLimit(),
|
||||
messageProperties.getWebSocketBufferSizeLimit()
|
||||
)
|
||||
);
|
||||
log.info("[connect] sessionId: {}, userId:{}, token:***{}", session.getId(), loginUser.getUserId(), StringUtils.right(token, 8));
|
||||
}
|
||||
|
||||
/**
|
||||
* 处理客户端发送的文本消息
|
||||
* 支持心跳ping/pong,以及自定义消息转发
|
||||
*
|
||||
* @param session WebSocket会话
|
||||
* @param message 文本消息
|
||||
*/
|
||||
@Override
|
||||
protected void handleTextMessage(WebSocketSession session, TextMessage message) {
|
||||
LoginUser loginUser = (LoginUser) session.getAttributes().get(MessageConstants.LOGIN_USER_KEY);
|
||||
if (ObjectUtil.isNull(loginUser)) {
|
||||
return;
|
||||
}
|
||||
|
||||
// 心跳处理:客户端发送ping,服务端回复pong
|
||||
if (MessageConstants.PING.equalsIgnoreCase(message.getPayload())) {
|
||||
webSocketSessionManager.sendMessage(session, MessageConstants.PONG);
|
||||
return;
|
||||
}
|
||||
|
||||
// 构建客户端自定义消息并发布
|
||||
PushDTO dto = PushDTO.of(List.of(loginUser.getUserId()), PushPayloadDTO.of(
|
||||
PushTypeEnum.CUSTOM,
|
||||
PushSourceEnum.CLIENT,
|
||||
message.getPayload(),
|
||||
null
|
||||
));
|
||||
webSocketSessionManager.publishMessage(dto);
|
||||
}
|
||||
|
||||
/**
|
||||
* 处理二进制消息(默认实现)
|
||||
*/
|
||||
@Override
|
||||
protected void handleBinaryMessage(WebSocketSession session, BinaryMessage message) throws Exception {
|
||||
super.handleBinaryMessage(session, message);
|
||||
}
|
||||
|
||||
/**
|
||||
* 处理Pong心跳响应
|
||||
* 维持长连接存活
|
||||
*/
|
||||
@Override
|
||||
protected void handlePongMessage(WebSocketSession session, PongMessage message) {
|
||||
webSocketSessionManager.sendPongMessage(session);
|
||||
}
|
||||
|
||||
/**
|
||||
* 传输异常处理
|
||||
* 记录异常日志
|
||||
*/
|
||||
@Override
|
||||
public void handleTransportError(WebSocketSession session, Throwable exception) {
|
||||
log.error("[transport error] sessionId: {}, exception:{}", session.getId(), exception.getMessage(), exception);
|
||||
LoginUser loginUser = (LoginUser) session.getAttributes().get(MessageConstants.LOGIN_USER_KEY);
|
||||
String token = (String) session.getAttributes().get(MessageConstants.LOGIN_TOKEN_KEY);
|
||||
if (ObjectUtil.hasNull(loginUser, token)) {
|
||||
webSocketSessionManager.closeSession(session, CloseStatus.SERVER_ERROR);
|
||||
return;
|
||||
}
|
||||
webSocketSessionManager.disconnect(loginUser.getUserId(), token, CloseStatus.SERVER_ERROR);
|
||||
}
|
||||
|
||||
/**
|
||||
* 连接关闭后触发
|
||||
* 注销用户会话
|
||||
*/
|
||||
@Override
|
||||
public void afterConnectionClosed(WebSocketSession session, CloseStatus status) {
|
||||
LoginUser loginUser = (LoginUser) session.getAttributes().get(MessageConstants.LOGIN_USER_KEY);
|
||||
String token = (String) session.getAttributes().get(MessageConstants.LOGIN_TOKEN_KEY);
|
||||
|
||||
if (ObjectUtil.hasNull(loginUser, token)) {
|
||||
log.info("[disconnect] invalid token received. sessionId: {}", session.getId());
|
||||
return;
|
||||
}
|
||||
|
||||
// 从会话管理器中移除连接
|
||||
webSocketSessionManager.disconnect(loginUser.getUserId(), token);
|
||||
log.info("[disconnect] sessionId: {}, userId:{}, token:***{}", session.getId(), loginUser.getUserId(), StringUtils.right(token, 8));
|
||||
}
|
||||
|
||||
/**
|
||||
* 是否支持分片消息
|
||||
* 关闭:不支持分片传输
|
||||
*/
|
||||
@Override
|
||||
public boolean supportsPartialMessages() {
|
||||
return false;
|
||||
}
|
||||
}
|
||||
+138
@@ -0,0 +1,138 @@
|
||||
package org.dromara.common.push.helper;
|
||||
|
||||
import lombok.AccessLevel;
|
||||
import lombok.NoArgsConstructor;
|
||||
import org.dromara.common.core.enums.PushSourceEnum;
|
||||
import org.dromara.common.core.enums.PushTypeEnum;
|
||||
import org.dromara.common.core.utils.SpringUtils;
|
||||
import org.dromara.common.push.core.PushSessionManager;
|
||||
import org.dromara.common.push.dto.PushDTO;
|
||||
import org.dromara.system.api.domain.PushPayloadDTO;
|
||||
|
||||
import java.util.List;
|
||||
|
||||
/**
|
||||
* 统一消息推送工具。
|
||||
*
|
||||
* @author Lion Li
|
||||
*/
|
||||
@NoArgsConstructor(access = AccessLevel.PRIVATE)
|
||||
public class PushHelper {
|
||||
|
||||
/**
|
||||
* 发送指定用户文本消息
|
||||
*
|
||||
* @param userId 目标用户ID
|
||||
* @param message 文本消息内容
|
||||
*/
|
||||
public static void sendMessage(Long userId, String message) {
|
||||
sendMessage(userId, buildMessage(message));
|
||||
}
|
||||
|
||||
/**
|
||||
* 全局广播文本消息
|
||||
*
|
||||
* @param message 文本消息内容
|
||||
*/
|
||||
public static void sendMessage(String message) {
|
||||
sendMessage(buildMessage(message));
|
||||
}
|
||||
|
||||
/**
|
||||
* 发送指定用户自定义消息体
|
||||
*
|
||||
* @param userId 目标用户ID
|
||||
* @param payload 消息推送体
|
||||
*/
|
||||
public static void sendMessage(Long userId, PushPayloadDTO payload) {
|
||||
if (!isEnabled()) {
|
||||
return;
|
||||
}
|
||||
getSessionManager().sendMessage(userId, payload);
|
||||
}
|
||||
|
||||
/**
|
||||
* 全局广播自定义消息体
|
||||
*
|
||||
* @param payload 消息推送体
|
||||
*/
|
||||
public static void sendMessage(PushPayloadDTO payload) {
|
||||
if (!isEnabled()) {
|
||||
return;
|
||||
}
|
||||
getSessionManager().sendMessage(payload);
|
||||
}
|
||||
|
||||
/**
|
||||
* 批量发布消息给指定用户列表
|
||||
*
|
||||
* @param userIds 用户ID集合
|
||||
* @param payload 消息推送体
|
||||
*/
|
||||
public static void publishMessage(List<Long> userIds, PushPayloadDTO payload) {
|
||||
publishMessage(PushDTO.of(userIds, payload));
|
||||
}
|
||||
|
||||
/**
|
||||
* 批量发布消息(使用完整推送DTO)
|
||||
*
|
||||
* @param dto 推送参数封装对象
|
||||
*/
|
||||
public static void publishMessage(PushDTO dto) {
|
||||
if (!isEnabled() || dto == null || dto.getPayload() == null) {
|
||||
return;
|
||||
}
|
||||
getSessionManager().publishMessage(dto);
|
||||
}
|
||||
|
||||
/**
|
||||
* 发布全局广播文本消息
|
||||
*
|
||||
* @param message 文本消息内容
|
||||
*/
|
||||
public static void publishAll(String message) {
|
||||
publishAll(buildMessage(message));
|
||||
}
|
||||
|
||||
/**
|
||||
* 发布全局广播自定义消息体
|
||||
*
|
||||
* @param payload 消息推送体
|
||||
*/
|
||||
public static void publishAll(PushPayloadDTO payload) {
|
||||
if (!isEnabled()) {
|
||||
return;
|
||||
}
|
||||
getSessionManager().publishAll(payload);
|
||||
}
|
||||
|
||||
/**
|
||||
* 判断消息推送功能是否开启
|
||||
* 读取配置:message.enabled
|
||||
*
|
||||
* @return 是否开启推送
|
||||
*/
|
||||
public static boolean isEnabled() {
|
||||
return Boolean.TRUE.equals(SpringUtils.getProperty("message.enabled", Boolean.class, Boolean.TRUE));
|
||||
}
|
||||
|
||||
/**
|
||||
* 获取推送会话管理器Bean
|
||||
*
|
||||
* @return PushSessionManager 实例
|
||||
*/
|
||||
private static PushSessionManager getSessionManager() {
|
||||
return SpringUtils.getBean(PushSessionManager.class);
|
||||
}
|
||||
|
||||
/**
|
||||
* 构建默认格式的消息推送体
|
||||
*
|
||||
* @param message 消息内容
|
||||
* @return 封装好的 PushPayloadDTO
|
||||
*/
|
||||
private static PushPayloadDTO buildMessage(String message) {
|
||||
return PushPayloadDTO.of(PushTypeEnum.MESSAGE, PushSourceEnum.BACKEND, message, null);
|
||||
}
|
||||
|
||||
}
|
||||
+44
@@ -0,0 +1,44 @@
|
||||
package org.dromara.common.push.interceptor;
|
||||
|
||||
import cn.dev33.satoken.stp.StpUtil;
|
||||
import org.dromara.common.push.constant.MessageConstants;
|
||||
import org.dromara.common.satoken.utils.LoginHelper;
|
||||
import org.dromara.system.api.model.LoginUser;
|
||||
import org.springframework.http.server.ServerHttpRequest;
|
||||
import org.springframework.http.server.ServerHttpResponse;
|
||||
import org.springframework.web.socket.WebSocketHandler;
|
||||
import org.springframework.web.socket.server.HandshakeInterceptor;
|
||||
|
||||
import java.util.Map;
|
||||
|
||||
/**
|
||||
* WebSocket 握手拦截器。
|
||||
*
|
||||
* @author Lion Li
|
||||
*/
|
||||
public class PlusWebSocketInterceptor implements HandshakeInterceptor {
|
||||
|
||||
/**
|
||||
* 握手前提取统一鉴权后的用户信息。
|
||||
*
|
||||
* @param attributes 用于传递到 WebSocketSession 的属性集合
|
||||
* @return 是否允许握手(true=允许,false=拒绝)
|
||||
*/
|
||||
@Override
|
||||
public boolean beforeHandshake(ServerHttpRequest request, ServerHttpResponse response, WebSocketHandler wsHandler,
|
||||
Map<String, Object> attributes) {
|
||||
LoginUser loginUser = LoginHelper.getLoginUser();
|
||||
String tokenValue = StpUtil.getTokenValue();
|
||||
attributes.put(MessageConstants.LOGIN_USER_KEY, loginUser);
|
||||
attributes.put(MessageConstants.LOGIN_TOKEN_KEY, tokenValue);
|
||||
return true;
|
||||
}
|
||||
|
||||
/**
|
||||
* 握手完成后触发
|
||||
* 此处无需处理,留空即可
|
||||
*/
|
||||
@Override
|
||||
public void afterHandshake(ServerHttpRequest request, ServerHttpResponse response, WebSocketHandler wsHandler, Exception exception) {
|
||||
}
|
||||
}
|
||||
+61
@@ -0,0 +1,61 @@
|
||||
package org.dromara.common.push.listener;
|
||||
|
||||
import cn.hutool.core.collection.CollUtil;
|
||||
import lombok.RequiredArgsConstructor;
|
||||
import lombok.extern.slf4j.Slf4j;
|
||||
import org.dromara.common.push.core.PushSessionManager;
|
||||
import org.springframework.boot.ApplicationArguments;
|
||||
import org.springframework.boot.ApplicationRunner;
|
||||
import org.springframework.core.Ordered;
|
||||
|
||||
/**
|
||||
* 统一消息主题订阅监听器。
|
||||
*
|
||||
* @author Lion Li
|
||||
*/
|
||||
@Slf4j
|
||||
@RequiredArgsConstructor
|
||||
public class MessageTopicListener implements ApplicationRunner, Ordered {
|
||||
|
||||
/**
|
||||
* 推送会话管理器
|
||||
*/
|
||||
private final PushSessionManager pushSessionManager;
|
||||
|
||||
/**
|
||||
* 项目启动后执行
|
||||
* 注册消息订阅,监听消息并分发给对应用户/全局广播
|
||||
*
|
||||
* @param args 启动参数
|
||||
*/
|
||||
@Override
|
||||
public void run(ApplicationArguments args) {
|
||||
// 订阅消息主题,处理消息分发
|
||||
pushSessionManager.subscribeMessage(message -> {
|
||||
if (message == null || message.getPayload() == null) {
|
||||
return;
|
||||
}
|
||||
log.info("消息主题订阅收到消息userIds={} message={}",
|
||||
message.getUserIds(),
|
||||
message.getPayload().getMessage());
|
||||
// 有指定用户 -> 单发
|
||||
if (CollUtil.isNotEmpty(message.getUserIds())) {
|
||||
message.getUserIds().forEach(userId -> pushSessionManager.sendMessage(userId, message.getPayload()));
|
||||
} else {
|
||||
// 无指定用户 -> 全局广播
|
||||
pushSessionManager.sendMessage(message.getPayload());
|
||||
}
|
||||
});
|
||||
log.info("初始化消息主题订阅监听器成功");
|
||||
}
|
||||
|
||||
/**
|
||||
* 执行顺序,优先级设为最高,确保消息订阅最先初始化
|
||||
*
|
||||
* @return 优先级,值越小越先执行
|
||||
*/
|
||||
@Override
|
||||
public int getOrder() {
|
||||
return -1;
|
||||
}
|
||||
}
|
||||
+55
@@ -0,0 +1,55 @@
|
||||
package org.dromara.common.push.properties;
|
||||
|
||||
import lombok.Data;
|
||||
import org.dromara.common.push.enums.MessageTransportEnum;
|
||||
import org.springframework.boot.context.properties.ConfigurationProperties;
|
||||
|
||||
/**
|
||||
* 统一消息推送配置。
|
||||
*
|
||||
* @author Lion Li
|
||||
*/
|
||||
@Data
|
||||
@ConfigurationProperties("message")
|
||||
public class MessageProperties {
|
||||
|
||||
/**
|
||||
* 是否启用消息推送。
|
||||
*/
|
||||
private Boolean enabled = true;
|
||||
|
||||
/**
|
||||
* 传输方式:sse / websocket。
|
||||
*/
|
||||
private String transport = MessageTransportEnum.SSE.getCode();
|
||||
|
||||
/**
|
||||
* 统一访问路径。
|
||||
*/
|
||||
private String path = "/resource/message";
|
||||
|
||||
/**
|
||||
* WebSocket 允许的跨域来源。
|
||||
*/
|
||||
private String[] allowedOrigins = {"*"};
|
||||
|
||||
/**
|
||||
* SSE 连接超时时间,单位毫秒。
|
||||
*/
|
||||
private long sseTimeout = 86_400_000L;
|
||||
|
||||
/**
|
||||
* 本地连接心跳检测间隔,单位秒。
|
||||
*/
|
||||
private long heartbeatInterval = 60L;
|
||||
|
||||
/**
|
||||
* WebSocket 单次发送超时时间,单位毫秒。
|
||||
*/
|
||||
private int webSocketSendTimeLimit = 10_000;
|
||||
|
||||
/**
|
||||
* WebSocket 发送缓冲区大小。
|
||||
*/
|
||||
private int webSocketBufferSizeLimit = 64_000;
|
||||
}
|
||||
+3
@@ -0,0 +1,3 @@
|
||||
org.dromara.common.push.config.MessageAutoConfiguration
|
||||
org.dromara.common.push.config.MessageSseConfiguration
|
||||
org.dromara.common.push.config.MessageWebSocketConfiguration
|
||||
Reference in New Issue
Block a user