| | |
| | | import java.util.Map; |
| | | import java.util.concurrent.ConcurrentHashMap; |
| | | import java.util.concurrent.ConcurrentMap; |
| | | import java.util.concurrent.Executor; |
| | | import java.util.regex.Matcher; |
| | | import java.util.regex.Pattern; |
| | | |
| | |
| | | @Slf4j |
| | | public final class PTZWebSocketUtils { |
| | | |
| | | // 存储 websocket session |
| | | public static final ConcurrentMap<String, WebSocketSession> ONLINE_USER_SESSIONS = new ConcurrentHashMap<>(); |
| | | |
| | | /** |
| | | * @param session 用户 session |
| | | * @param message 发送内容 |
| | | */ |
| | | /** 推送线程池,由 PTZWebSocketInitializer 在启动时注入;未注入时降级为同步执行 */ |
| | | private static volatile Executor pushExecutor; |
| | | |
| | | public static void setPushExecutor(Executor executor) { |
| | | pushExecutor = executor; |
| | | } |
| | | |
| | | public static void sendMessage(WebSocketSession session, String message) { |
| | | if (session == null) { |
| | | return; |
| | |
| | | } |
| | | synchronized (session) { |
| | | try { |
| | | //log.debug("发送消息:" + message); |
| | | session.sendMessage(new TextMessage(String.join(", ", message))); |
| | | session.sendMessage(new TextMessage(message)); |
| | | } catch (IOException e) { |
| | | log.error("sendMessage IOException ", e); |
| | | } |
| | | } |
| | | } |
| | | |
| | | /** |
| | | * @param session 用户 session |
| | | * @param message 发送内容 |
| | | */ |
| | | public static void sendMessage(WebSocketSession session, Map message) { |
| | | if (session == null) { |
| | | return; |
| | | } |
| | | final InetSocketAddress remoteAddress = session.getRemoteAddress(); |
| | | if (remoteAddress == null) { |
| | | return; |
| | | } |
| | | synchronized (session) { |
| | | try { |
| | | session.sendMessage(new TextMessage(JSONObject.toJSONString(message))); |
| | | } catch (IOException e) { |
| | | log.error("sendMessage IOException ", e); |
| | | } |
| | | } |
| | | sendMessage(session, JSONObject.toJSONString(message)); |
| | | } |
| | | |
| | | public static void sendMessage(WebSocketSession session, List message) { |
| | | if (session == null) { |
| | | return; |
| | | } |
| | | final InetSocketAddress remoteAddress = session.getRemoteAddress(); |
| | | if (remoteAddress == null) { |
| | | return; |
| | | } |
| | | synchronized (session) { |
| | | try { |
| | | session.sendMessage(new TextMessage(JSONObject.toJSONString(message))); |
| | | } catch (IOException e) { |
| | | log.error("sendMessage IOException ", e); |
| | | } |
| | | } |
| | | sendMessage(session, JSON.toJSONString(message)); |
| | | } |
| | | |
| | | /** |
| | | * 推送消息到其他客户端 |
| | | * |
| | | * @param message |
| | | * 推送字符串消息到所有客户端 |
| | | */ |
| | | public static void sendMessageAll(String message) { |
| | | ONLINE_USER_SESSIONS.forEach((sessionId, session) -> sendMessage(session, message)); |
| | | broadcast(message); |
| | | } |
| | | |
| | | /** |
| | | * 推送消息到其他客户端 |
| | | * |
| | | * @param message |
| | | * 推送Map消息到所有客户端 |
| | | */ |
| | | public static void sendMessageAll(Map message) { |
| | | JSONObject jsonObject = new JSONObject(message); |
| | | ONLINE_USER_SESSIONS.forEach((sessionId, session) -> sendMessage(session, jsonObject.toString())); |
| | | broadcast(JSONObject.toJSONString(message)); |
| | | } |
| | | |
| | | /** |
| | | * 推送消息到其他客户端 |
| | | * |
| | | * @param message |
| | | * 推送List消息到所有客户端 |
| | | */ |
| | | public static void sendMessageAll(List message) { |
| | | ONLINE_USER_SESSIONS.forEach((sessionId, session) -> sendMessage(session, JSON.toJSONString(message))); |
| | | broadcast(JSON.toJSONString(message)); |
| | | } |
| | | |
| | | /** |
| | | * @Author 刘苏义 |
| | | * @Description 发送消息给userId_前缀的人 |
| | | * @Date 2024/7/16 10:24 |
| | | * @Param |
| | | * @return |
| | | * 一次序列化,并行分发到所有 session |
| | | */ |
| | | private static void broadcast(String payload) { |
| | | if (ONLINE_USER_SESSIONS.isEmpty()) return; |
| | | Executor executor = pushExecutor; |
| | | if (executor == null) { |
| | | ONLINE_USER_SESSIONS.values().forEach(session -> sendMessage(session, payload)); |
| | | return; |
| | | } |
| | | ONLINE_USER_SESSIONS.values().forEach(session -> |
| | | executor.execute(() -> sendMessage(session, payload))); |
| | | } |
| | | |
| | | /** |
| | | * 发送消息给 userId_ 前缀的人 |
| | | */ |
| | | public static void sendMessagePrefix(String targetId, String message) { |
| | | String regex = "^" + Pattern.quote(targetId) + "_\\d+$"; |
| | |
| | | return matcher.matches(); |
| | | }).map(Map.Entry::getValue).forEach(session -> PTZWebSocketUtils.sendMessage(session, message)); |
| | | } |
| | | } |
| | | } |