package com.maibu.websocket; import com.maibu.common.CommandConstant; import com.maibu.common.MiddleConstant; import com.maibu.dto.*; import com.maibu.enums.CommandRequestType; import com.maibu.memory.MiddleGlobalMemory; import com.maibu.netty.NettyClient; import com.maibu.utils.CommandUtils; import com.maibu.utils.ControlUtils; import com.maibu.utils.json.JsonUtils; import io.netty.buffer.ByteBuf; import io.netty.buffer.Unpooled; import org.apache.commons.lang3.StringUtils; import org.java_websocket.WebSocket; import org.java_websocket.handshake.ClientHandshake; import org.java_websocket.server.WebSocketServer; import org.slf4j.Logger; import org.slf4j.LoggerFactory; import java.net.InetSocketAddress; import java.util.Collection; import java.util.List; import java.util.Map; import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.Executors; import java.util.concurrent.ScheduledExecutorService; import java.util.concurrent.TimeUnit; /** * WebSocket服务器实现类 * 使用java-websocket库实现WebSocket服务端功能 */ public class WebsocketHandler extends WebSocketServer { private static final Logger logger = LoggerFactory.getLogger(WebsocketHandler.class); // 存储客户端连接:key=WebSocket连接对象,value=最后一次活跃时间(毫秒) private final Map clientActiveTimeMap = new ConcurrentHashMap<>(); // 空闲超时时间(10秒):超过10秒无消息则判定为断开 private static final long IDLE_TIMEOUT = 10 * 1000; // 定时检测线程池(每隔10秒检测一次空闲连接) private final ScheduledExecutorService idleCheckExecutor = Executors.newSingleThreadScheduledExecutor(); public WebsocketHandler(int port) { super(new InetSocketAddress(port)); // 启动空闲连接检测任务 // startIdleCheckTask(); } @Override public void onOpen(WebSocket conn, ClientHandshake handshake) { logger.info("客户端连接成功: remoteAddress={}, 当前在线数: {}", conn.getRemoteSocketAddress(), MiddleGlobalMemory.onlineSockets.size()); } @Override public void onClose(WebSocket conn, int code, String reason, boolean remote) { String key = MiddleGlobalMemory.findBySocket(conn); if (key != null) { releaseResource(key); logger.info("客户端断开连接: remoteAddress={}, code={}, reason={}, 当前在线数: {}", conn.getRemoteSocketAddress(), code, reason, MiddleGlobalMemory.onlineSockets.size()); } } @Override public void onMessage(WebSocket conn, String message) { try { WebMessageBaseDTO webMessageBaseDTO = JsonUtils.parseObject(message, WebMessageBaseDTO.class); logger.debug("收到websocket消息:{}", message); if (webMessageBaseDTO != null) { String token = webMessageBaseDTO.getToken(); String userName = webMessageBaseDTO.getUserName(); String deviceId = webMessageBaseDTO.getDeviceId(); MesType type = webMessageBaseDTO.getType(); String key = userName + ":web" + ":" + token; if (MesType.AUTH.equals(type)) { if (StringUtils.isEmpty(deviceId)) { return; } //todo 先关闭原先的 WebSocket oldWebsocket = MiddleGlobalMemory.onlineSockets.get(key); NettyClient oldNettyClient = MiddleGlobalMemory.nettyClientMap.get(key); if (oldWebsocket != null && oldWebsocket.isOpen()) { oldWebsocket.close(); } if (oldNettyClient != null && oldNettyClient.isConnected()) { oldNettyClient.disconnect(); } MiddleGlobalMemory.onlineSockets.put(key, conn); //这里可以根据需要初始化Netty客户端 NettyClient client = new NettyClient(); client.key = key; client.requestDeviceId = deviceId; client.token = token; client.userName = userName; client.initClient(); client.initConnect(MiddleConstant.vertxIp, 9001); MiddleGlobalMemory.nettyClientMap.put(key, client); } else if (MesType.path.equals(type)) { //todo下发路径 预留 pathDistribute(); } else if (MesType.remoteControl.equals(type)) { //遥控 RemoteData remoteData = webMessageBaseDTO.getRemoteData(); if (remoteData != null) { List axes = remoteData.getAxes(); List buttons = remoteData.getButtons(); Double y = axes.get(1) * 1000 * -1; //前后值反了 Double x = axes.get(2) * 1000; int intX = (int) x.doubleValue(); int intY = (int) y.doubleValue(); MiddleGlobalMemory.initDiff(); LawnMoverProtocolEntity entity = ControlUtils.getCarControlMessage(MiddleGlobalMemory.diffSteer, intX, intY, buttons); byte[] buf = ControlUtils.buildControlCommand(entity); logger.debug("Netty客户端:{}发送遥控数据: {}", userName, JsonUtils.toJsonString(entity)); NettyClient nettyClient = MiddleGlobalMemory.nettyClientMap.get(key); if (nettyClient != null && nettyClient.isConnected()) { // if (!nettyClient.isConnected()) nettyClient.connect(); nettyClient.sendData(buf); } } } // else if(MesType.heartbeat.equals(type)){ // clientActiveTimeMap.put(conn,System.currentTimeMillis()); // } else if(MesType.switchPermission.equals(type)){ DeviceRequestDTO requestDTO = new DeviceRequestDTO(); requestDTO.setRequest(CommandRequestType.switch_control); requestDTO.setToken(token); requestDTO.setPlatform("web"); requestDTO.setDeviceId(deviceId); requestDTO.setUserId(userName); logger.debug("Netty客户端:{}发送操控请求数据: {}", userName, JsonUtils.toJsonString(requestDTO)); NettyClient nettyClient = MiddleGlobalMemory.nettyClientMap.get(key); ByteBuf buf = Unpooled.buffer(); CommandUtils.buildCommand(buf,requestDTO, CommandConstant.interaction); if (nettyClient != null && nettyClient.isConnected()) { // if (!nettyClient.isConnected()) nettyClient.connect(); nettyClient.sendData(buf.array()); } } else if(MesType.switchResult.equals(type)){ DeviceRequestDTO requestDTO = new DeviceRequestDTO(); requestDTO.setToken(token); requestDTO.setPlatform("web"); requestDTO.setDeviceId(deviceId); requestDTO.setUserId(userName); DeviceRespondDTO respond = new DeviceRespondDTO(); respond.setSwitchResult(webMessageBaseDTO.getSwitchResult()); respond.setDeviceId(deviceId); requestDTO.setRespond(respond); logger.debug("Netty客户端:{}发送切换操控权请求反馈: {}", userName, JsonUtils.toJsonString(requestDTO)); NettyClient nettyClient = MiddleGlobalMemory.nettyClientMap.get(key); ByteBuf buf = Unpooled.buffer(); CommandUtils.buildCommand(buf,requestDTO, CommandConstant.interaction); if (nettyClient != null && nettyClient.isConnected()) { // if (!nettyClient.isConnected()) nettyClient.connect(); nettyClient.sendData(buf.array()); } } } } catch (Exception e) { logger.error("onMessage error:{}", e.getMessage()); } } public void releaseResource(String key) { MiddleGlobalMemory.onlineSockets.remove(key); if (MiddleGlobalMemory.nettyClientMap.containsKey(key)) { NettyClient client = MiddleGlobalMemory.nettyClientMap.get(key); if (client != null) { if (client.isConnected()) { client.disconnect(); } MiddleGlobalMemory.nettyClientMap.remove(key); } } } public void pathDistribute() { } @Override public void onError(WebSocket conn, Exception ex) { String key = MiddleGlobalMemory.findBySocket(conn); logger.error("WebSocket错误: clientId={}, remoteAddress={}, error={}", key, conn != null ? conn.getRemoteSocketAddress() : null, ex.getMessage(), ex); // 发生错误时清理连接 if (key != null) { releaseResource(key); } } @Override public void onStart() { logger.info("WebSocket服务端启动成功,监听端口: {}", getPort()); } /** * 给指定WebSocket连接发送消息 */ public static void sendMessageToClient(WebSocket conn, String message) { if (conn != null && conn.isOpen()) { conn.send(message); String clientId = MiddleGlobalMemory.findBySocket(conn); logger.debug("发送消息给客户端: clientId={}, message={}", clientId, message); } } /** * 广播消息给所有连接的客户端 */ public void broadcastMessage(String message) { Collection connections = getConnections(); for (WebSocket conn : connections) { if (conn.isOpen()) { conn.send(message); } } logger.debug("广播消息给所有客户端: message={}, 客户端数量: {}", message, connections.size()); } // 核心:定时检测空闲连接(每隔10秒执行一次) private void startIdleCheckTask() { idleCheckExecutor.scheduleAtFixedRate(() -> { long now = System.currentTimeMillis(); // 遍历所有连接,检测是否超时 for (Map.Entry entry : clientActiveTimeMap.entrySet()) { WebSocket conn = entry.getKey(); Long lastActiveTime = entry.getValue(); // 超时判定:当前时间 - 最后活跃时间 > 空闲超时时间 if (now - lastActiveTime > IDLE_TIMEOUT) { System.out.println("客户端空闲超时:" + conn.getRemoteSocketAddress() + ",强制关闭连接"); // 强制关闭连接(会触发 onClose,但 remote=true 表示服务端主动关闭) conn.close(1001, "idle timeout"); } } }, 10, 10, TimeUnit.SECONDS); } // 核心新增:停止WebSocket服务的方法(释放端口+线程池+Netty资源) public void stopServer() throws InterruptedException { logger.info("开始关闭WebSocket服务,监听端口: {}", getPort()); // 1. 关闭所有客户端连接 Collection connections = getConnections(); for (WebSocket conn : connections) { if (conn.isOpen()) { try { conn.close(1000, "server shutdown"); // 正常关闭客户端 String key = MiddleGlobalMemory.findBySocket(conn); if (key != null) { releaseResource(key); // 清理对应Netty资源 } } catch (Exception e) { logger.warn("关闭客户端连接失败: {}", conn.getRemoteSocketAddress(), e); } } } // clientActiveTimeMap.clear(); // 清空空闲记录 // // // 2. 关闭定时检测线程池(核心:终止非守护线程) // if (!idleCheckExecutor.isShutdown()) { // idleCheckExecutor.shutdownNow(); // 立即关闭线程池 // try { // // 等待线程池终止,最多等3秒 // if (!idleCheckExecutor.awaitTermination(3, TimeUnit.SECONDS)) { // logger.warn("空闲检测线程池未正常终止"); // } // } catch (InterruptedException e) { // idleCheckExecutor.shutdownNow(); // Thread.currentThread().interrupt(); // } // } // 3. 批量关闭所有Netty客户端连接 for (NettyClient client : MiddleGlobalMemory.nettyClientMap.values()) { if (client != null && client.isConnected()) { try { client.disconnect(); } catch (Exception e) { logger.warn("关闭Netty客户端失败", e); } } } MiddleGlobalMemory.nettyClientMap.clear(); MiddleGlobalMemory.onlineSockets.clear(); // 4. 关闭WebSocketServer本身,释放9002端口(java-websocket核心方法) try { this.stop(3000); // 3秒超时关闭服务端 logger.info("WebSocket服务端关闭成功,端口: {}", getPort()); } catch (InterruptedException e) { logger.error("关闭WebSocket服务端失败", e); this.stop(); // 强制关闭 } } }