通过WebSocket发送控制指令

This commit is contained in:
shan 2026-01-24 12:12:55 +08:00
parent 85d0472f91
commit dd1d408688
3 changed files with 137 additions and 4 deletions

View File

@ -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");
}
}

View File

@ -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://<host>: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<String, WebSocketSession> sessions = new ConcurrentHashMap<>();
private final Map<String, ScheduledFuture<?>> 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());
}
}
}

View File

@ -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()