diff --git a/qaup-collision/src/main/java/com/qaup/collision/websocket/config/WebSocketConfig.java b/qaup-collision/src/main/java/com/qaup/collision/websocket/config/WebSocketConfig.java index bbe2eff..ff07a4c 100644 --- a/qaup-collision/src/main/java/com/qaup/collision/websocket/config/WebSocketConfig.java +++ b/qaup-collision/src/main/java/com/qaup/collision/websocket/config/WebSocketConfig.java @@ -7,6 +7,7 @@ import org.springframework.web.socket.config.annotation.WebSocketConfigurer; import org.springframework.web.socket.config.annotation.WebSocketHandlerRegistry; import com.qaup.collision.websocket.handler.CollisionWebSocketHandler; +import com.qaup.collision.websocket.handler.VehicleCommandInfoWebSocketHandler; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -24,9 +25,12 @@ public class WebSocketConfig implements WebSocketConfigurer { private static final Logger logger = LoggerFactory.getLogger(WebSocketConfig.class); private final CollisionWebSocketHandler collisionWebSocketHandler; + private final VehicleCommandInfoWebSocketHandler vehicleCommandInfoWebSocketHandler; - public WebSocketConfig(CollisionWebSocketHandler collisionWebSocketHandler) { + public WebSocketConfig(CollisionWebSocketHandler collisionWebSocketHandler, + VehicleCommandInfoWebSocketHandler vehicleCommandInfoWebSocketHandler) { this.collisionWebSocketHandler = collisionWebSocketHandler; + this.vehicleCommandInfoWebSocketHandler = vehicleCommandInfoWebSocketHandler; logger.info("🚀 WebSocket配置类初始化..."); } @@ -37,10 +41,14 @@ public class WebSocketConfig implements WebSocketConfigurer { // 注册冲突检测WebSocket端点 registry.addHandler(collisionWebSocketHandler, "/collision") .setAllowedOrigins("*"); // 允许所有来源,生产环境应该限制具体域名 + + // 注册车辆控制指令测试端点 + registry.addHandler(vehicleCommandInfoWebSocketHandler, "/VehicleCommandInfo") + .setAllowedOrigins("*"); logger.info("✅ WebSocket端点注册完成"); - logger.info("🎯 端点路径: /collision"); + logger.info("🎯 端点路径: /collision, /VehicleCommandInfo"); logger.info("🌐 允许的来源: *"); - logger.info("📡 WebSocket服务可用: ws://localhost:8080/collision"); + logger.info("📡 WebSocket服务可用: ws://localhost:8080/collision, ws://localhost:8080/VehicleCommandInfo"); } } \ No newline at end of file diff --git a/qaup-collision/src/main/java/com/qaup/collision/websocket/handler/VehicleCommandInfoWebSocketHandler.java b/qaup-collision/src/main/java/com/qaup/collision/websocket/handler/VehicleCommandInfoWebSocketHandler.java new file mode 100644 index 0000000..bfc427c --- /dev/null +++ b/qaup-collision/src/main/java/com/qaup/collision/websocket/handler/VehicleCommandInfoWebSocketHandler.java @@ -0,0 +1,125 @@ +package com.qaup.collision.websocket.handler; + +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; +import org.springframework.lang.NonNull; +import org.springframework.stereotype.Component; +import org.springframework.web.socket.CloseStatus; +import org.springframework.web.socket.TextMessage; +import org.springframework.web.socket.WebSocketHandler; +import org.springframework.web.socket.WebSocketMessage; +import org.springframework.web.socket.WebSocketSession; + +import java.io.IOException; +import java.util.Map; +import java.util.concurrent.ConcurrentHashMap; +import java.util.concurrent.Executors; +import java.util.concurrent.ScheduledExecutorService; +import java.util.concurrent.ScheduledFuture; +import java.util.concurrent.TimeUnit; + +/** + * WebSocket 车辆控制指令测试端点 + * + * 客户端可通过 ws://:8080/VehicleCommandInfo 连接, + * 服务端会周期性向所有连接会话发送固定的车辆控制指令 JSON。 + */ +@Component +public class VehicleCommandInfoWebSocketHandler implements WebSocketHandler { + + private static final Logger LOGGER = LoggerFactory.getLogger(VehicleCommandInfoWebSocketHandler.class); + + private static final String COMMAND_JSON = "{" + + "\"messageUniqueId\": \"68f79d1a-e27f-11ed-b28c-2cf05d9c2649\"," + + "\"timestamp\": 1736175610," + + "\"vehicleID\": \"A001\"," + + "\"commandType\": \"SIGNAL\"," + + "\"commandReason\": \"TRAFFIC_LIGHT\"," + + "\"signalState\":\"RED\"," + + "\"intersectionId\":\"002\"," + + "\"latitude\": 343.23," + + "\"longitude\": 343.23," + + "\"relativeSpeed\": 3," + + "\"relativeMotionX\": 2002.12," + + "\"relativeMotionY\":100.12," + + "\"minDistance\":10.5" + + "}"; + + private final Map sessions = new ConcurrentHashMap<>(); + private final Map> sessionTasks = new ConcurrentHashMap<>(); + private final ScheduledExecutorService scheduler = Executors.newScheduledThreadPool(1); + + @Override + public void afterConnectionEstablished(@NonNull WebSocketSession session) { + String sessionId = session.getId(); + sessions.put(sessionId, session); + LOGGER.info("VehicleCommandInfo WebSocket 连接建立, sessionId={}", sessionId); + + ScheduledFuture task = scheduler.scheduleAtFixedRate( + () -> sendCommandJson(session), + 0, + 1, + TimeUnit.SECONDS + ); + sessionTasks.put(sessionId, task); + } + + @Override + public void handleMessage(@NonNull WebSocketSession session, @NonNull WebSocketMessage message) { + // 简单回显客户端消息,主要用于连通性测试 + LOGGER.debug("VehicleCommandInfo 收到客户端消息, sessionId={}, payload={}", + session.getId(), message.getPayload()); + } + + @Override + public void handleTransportError(@NonNull WebSocketSession session, @NonNull Throwable exception) { + LOGGER.warn("VehicleCommandInfo WebSocket 传输错误, sessionId={}", session.getId(), exception); + cleanupSession(session, CloseStatus.SERVER_ERROR); + } + + @Override + public void afterConnectionClosed(@NonNull WebSocketSession session, @NonNull CloseStatus closeStatus) { + LOGGER.info("VehicleCommandInfo WebSocket 连接关闭, sessionId={}, status={}", session.getId(), closeStatus); + cleanupSession(session, closeStatus); + } + + @Override + public boolean supportsPartialMessages() { + return false; + } + + private void sendCommandJson(WebSocketSession session) { + if (session == null || !session.isOpen()) { + return; + } + try { + synchronized (session) { + if (session.isOpen()) { + session.sendMessage(new TextMessage(COMMAND_JSON)); + } + } + } catch (IOException e) { + LOGGER.warn("向 sessionId={} 发送车辆控制指令 JSON 失败: {}", session.getId(), e.getMessage()); + cleanupSession(session, CloseStatus.SESSION_NOT_RELIABLE); + } catch (Exception e) { + LOGGER.warn("向 sessionId={} 发送车辆控制指令 JSON 时发生异常", session.getId(), e); + cleanupSession(session, CloseStatus.SERVER_ERROR); + } + } + + private void cleanupSession(WebSocketSession session, CloseStatus status) { + String sessionId = session.getId(); + ScheduledFuture task = sessionTasks.remove(sessionId); + if (task != null) { + task.cancel(true); + } + sessions.remove(sessionId); + try { + if (session.isOpen()) { + session.close(status); + } + } catch (IOException e) { + LOGGER.debug("关闭 sessionId={} 时出现异常: {}", sessionId, e.getMessage()); + } + } +} diff --git a/qaup-framework/src/main/java/com/qaup/framework/config/SecurityConfig.java b/qaup-framework/src/main/java/com/qaup/framework/config/SecurityConfig.java index ac6f029..62794f0 100644 --- a/qaup-framework/src/main/java/com/qaup/framework/config/SecurityConfig.java +++ b/qaup-framework/src/main/java/com/qaup/framework/config/SecurityConfig.java @@ -112,7 +112,7 @@ public class SecurityConfig // 对于登录login 注册register 验证码captchaImage 允许匿名访问 requests.requestMatchers("/login", "/register", "/captchaImage").permitAll() // WebSocket端点,允许匿名访问 - .requestMatchers("/collision", "/collision/**", "/test/websocket/**").permitAll() + .requestMatchers("/collision", "/collision/**", "/VehicleCommandInfo", "/test/websocket/**").permitAll() // 静态资源,可匿名访问 .requestMatchers(HttpMethod.GET, "/", "/*.html", "/**.html", "/**.css", "/**.js", "/profile/**").permitAll() .requestMatchers("/swagger-ui.html", "/v3/api-docs/**", "/swagger-ui/**", "/druid/**").permitAll()