feat(controlled): 实现设备激活与安全认证流程
- 添加API客户端、加密存储和provision/token激活逻辑 - WebSocket改用Bearer令牌认证,移除REGISTER请求 - 设备ID改为服务端下发,支持令牌刷新和强制下线处理 - 新增deviceSecret加密存储和accessToken自动刷新 - 更新设备ID获取方式为出厂SN,添加安全存储依赖
This commit is contained in:
@@ -4,8 +4,16 @@ import com.fasterxml.jackson.databind.ObjectMapper;
|
||||
import com.ttstd.signaling.manager.ConnectionRequestManager;
|
||||
import com.ttstd.signaling.manager.SessionManager;
|
||||
import com.ttstd.signaling.manager.SignalMetrics;
|
||||
import com.ttstd.signaling.model.AuthPrincipal;
|
||||
import com.ttstd.signaling.model.DeviceBinding;
|
||||
import com.ttstd.signaling.model.DeviceType;
|
||||
import com.ttstd.signaling.model.SignalMessage;
|
||||
import com.ttstd.signaling.security.AuthHandshakeInterceptor;
|
||||
import com.ttstd.signaling.service.AccountService;
|
||||
import com.ttstd.signaling.service.AuditService;
|
||||
import com.ttstd.signaling.service.BindingService;
|
||||
import com.ttstd.signaling.service.DeviceIdentityService;
|
||||
import jakarta.annotation.PostConstruct;
|
||||
import org.slf4j.Logger;
|
||||
import org.slf4j.LoggerFactory;
|
||||
import org.springframework.stereotype.Component;
|
||||
@@ -18,6 +26,7 @@ import java.io.IOException;
|
||||
import java.util.HashMap;
|
||||
import java.util.List;
|
||||
import java.util.Map;
|
||||
import java.util.Set;
|
||||
|
||||
@Component
|
||||
public class SignalWebSocketHandler extends TextWebSocketHandler {
|
||||
@@ -25,40 +34,129 @@ public class SignalWebSocketHandler extends TextWebSocketHandler {
|
||||
private static final Logger logger = LoggerFactory.getLogger(SignalWebSocketHandler.class);
|
||||
private static final ObjectMapper objectMapper = new ObjectMapper();
|
||||
|
||||
/** 因鉴权/封禁被关闭连接时使用的关闭码 */
|
||||
private static final CloseStatus CLOSE_UNAUTHORIZED = new CloseStatus(4001, "UNAUTHORIZED");
|
||||
private static final CloseStatus CLOSE_FORCED_LOGOUT = new CloseStatus(4003, "FORCED_LOGOUT");
|
||||
|
||||
private final SessionManager sessionManager;
|
||||
private final ConnectionRequestManager connectionRequestManager;
|
||||
private final SignalMetrics metrics;
|
||||
private final AccountService accountService;
|
||||
private final DeviceIdentityService deviceIdentityService;
|
||||
private final BindingService bindingService;
|
||||
private final AuditService auditService;
|
||||
|
||||
public SignalWebSocketHandler(SessionManager sessionManager,
|
||||
ConnectionRequestManager connectionRequestManager,
|
||||
SignalMetrics metrics) {
|
||||
SignalMetrics metrics,
|
||||
AccountService accountService,
|
||||
DeviceIdentityService deviceIdentityService,
|
||||
BindingService bindingService,
|
||||
AuditService auditService) {
|
||||
this.sessionManager = sessionManager;
|
||||
this.connectionRequestManager = connectionRequestManager;
|
||||
this.metrics = metrics;
|
||||
this.accountService = accountService;
|
||||
this.deviceIdentityService = deviceIdentityService;
|
||||
this.bindingService = bindingService;
|
||||
this.auditService = auditService;
|
||||
this.connectionRequestManager.setSender(this::sendToDevice);
|
||||
}
|
||||
|
||||
/**
|
||||
* 注册会话失效回调:账号被封禁 / 会话被踢出 / 设备被禁用时立即断开对应连接。
|
||||
*/
|
||||
@PostConstruct
|
||||
void registerRevocationListener() {
|
||||
AccountService.SessionRevocationListener listener = this::forceDisconnect;
|
||||
accountService.setRevocationListener(listener);
|
||||
deviceIdentityService.setRevocationListener(listener);
|
||||
}
|
||||
|
||||
/**
|
||||
* 强制断开指定主体的连接。
|
||||
*
|
||||
* @param principalId 主体 ID
|
||||
* @param sessionId 指定登录会话;为 null 表示断开该主体全部连接
|
||||
*/
|
||||
private void forceDisconnect(String principalId, String sessionId, String reason) {
|
||||
Set<String> deviceIds = sessionManager.getDeviceIdsByPrincipal(principalId);
|
||||
for (String deviceId : deviceIds) {
|
||||
AuthPrincipal principal = sessionManager.getPrincipal(deviceId);
|
||||
if (principal == null) {
|
||||
continue;
|
||||
}
|
||||
// 指定了会话时只断开该会话,避免误伤同账号其他端
|
||||
if (sessionId != null && !sessionId.equals(principal.sessionId())) {
|
||||
continue;
|
||||
}
|
||||
WebSocketSession session = sessionManager.getSession(deviceId);
|
||||
if (session == null || !session.isOpen()) {
|
||||
continue;
|
||||
}
|
||||
Map<String, Object> notice = new HashMap<>();
|
||||
notice.put("type", "FORCE_LOGOUT");
|
||||
notice.put("payload", reason == null ? "会话已失效" : reason);
|
||||
sendToSession(session, notice);
|
||||
try {
|
||||
session.close(CLOSE_FORCED_LOGOUT);
|
||||
} catch (IOException e) {
|
||||
logger.warn("强制断开连接 {} 失败: {}", deviceId, e.getMessage());
|
||||
}
|
||||
logger.info("已强制断开连接: deviceId={} principal={} reason={}",
|
||||
deviceId, principalId, reason);
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* 向指定设备回送服务端事件(错误/超时通知等)。
|
||||
*/
|
||||
private void sendToDevice(String deviceId, Object message) {
|
||||
WebSocketSession session = sessionManager.getSession(deviceId);
|
||||
if (session == null || !session.isOpen()) {
|
||||
logger.warn("Cannot send event to {}: device offline", deviceId);
|
||||
logger.warn("无法向 {} 发送事件:连接不在线", deviceId);
|
||||
return;
|
||||
}
|
||||
sendToSession(session, message);
|
||||
}
|
||||
|
||||
/**
|
||||
* 连接建立后立即以握手阶段裁定的身份完成注册,无需客户端再发 REGISTER。
|
||||
*/
|
||||
@Override
|
||||
public void afterConnectionEstablished(WebSocketSession session) throws Exception {
|
||||
logger.info("New WebSocket connection: {}", session.getId());
|
||||
AuthPrincipal principal = principalOf(session);
|
||||
if (principal == null) {
|
||||
// 正常情况下握手拦截器已拦截,此处为纵深防御
|
||||
logger.warn("连接 {} 缺少鉴权主体,立即关闭", session.getId());
|
||||
session.close(CLOSE_UNAUTHORIZED);
|
||||
return;
|
||||
}
|
||||
|
||||
sessionManager.registerDevice(principal, session);
|
||||
metrics.incRegistered();
|
||||
metrics.recordSessionCount(sessionManager.getOnlineCount());
|
||||
|
||||
Map<String, Object> response = new HashMap<>();
|
||||
response.put("type", "REGISTER_SUCCESS");
|
||||
response.put("deviceId", principal.deviceId());
|
||||
response.put("deviceType", principal.deviceType().name());
|
||||
response.put("displayName", principal.displayName());
|
||||
sendToSession(session, response);
|
||||
|
||||
logger.info("连接已建立并注册: {} ({})", principal.deviceId(), principal.deviceType());
|
||||
}
|
||||
|
||||
@Override
|
||||
protected void handleTextMessage(WebSocketSession session, TextMessage message) throws Exception {
|
||||
AuthPrincipal principal = principalOf(session);
|
||||
if (principal == null) {
|
||||
session.close(CLOSE_UNAUTHORIZED);
|
||||
return;
|
||||
}
|
||||
|
||||
String payload = message.getPayload();
|
||||
logger.debug("Received message: {}", payload);
|
||||
logger.debug("收到消息: {}", payload);
|
||||
metrics.incMessage();
|
||||
|
||||
try {
|
||||
@@ -66,22 +164,54 @@ public class SignalWebSocketHandler extends TextWebSocketHandler {
|
||||
String type = signalMessage.getType();
|
||||
|
||||
if (type == null) {
|
||||
logger.warn("Message type is null");
|
||||
logger.warn("消息缺少 type 字段");
|
||||
return;
|
||||
}
|
||||
|
||||
// 关键:发送方身份一律以服务端鉴权结果覆盖,忽略客户端自报值,
|
||||
// 防止伪造 fromDeviceId 冒充他人。
|
||||
String claimedFrom = signalMessage.getFromDeviceId();
|
||||
if (claimedFrom != null && !claimedFrom.equals(principal.deviceId())) {
|
||||
logger.warn("客户端自报 fromDeviceId={} 与鉴权身份 {} 不一致,已强制覆盖",
|
||||
claimedFrom, principal.deviceId());
|
||||
}
|
||||
signalMessage.setFromDeviceId(principal.deviceId());
|
||||
signalMessage.setDeviceType(principal.deviceType().name());
|
||||
|
||||
switch (type.toUpperCase()) {
|
||||
case "REGISTER":
|
||||
handleRegister(session, signalMessage);
|
||||
// 身份已在握手阶段确定,REGISTER 仅作兼容响应
|
||||
handleLegacyRegister(session, principal);
|
||||
break;
|
||||
case "DEVICE_LIST":
|
||||
handleDeviceList(session, signalMessage);
|
||||
handleDeviceList(session, principal);
|
||||
break;
|
||||
case "PING":
|
||||
// 客户端心跳保活消息,无需处理,仅用于防止中间代理因空闲超时断开连接
|
||||
break;
|
||||
case "OFFER":
|
||||
metrics.incOffer();
|
||||
// 绑定/黑名单前置校验:仅允许已绑定且未被拉黑的主控端发起
|
||||
if (principal.deviceType() == DeviceType.CONTROLLER) {
|
||||
String targetUid = signalMessage.getToDeviceId();
|
||||
if (targetUid == null || targetUid.isEmpty()) {
|
||||
sendError(session, "缺少目标设备ID");
|
||||
metrics.incBlockedOffer();
|
||||
break;
|
||||
}
|
||||
if (!bindingService.isBound(targetUid, principal.principalId())) {
|
||||
sendError(session, "未与该设备建立绑定关系,无法发起连接");
|
||||
metrics.incBlockedOffer();
|
||||
auditOfferBlocked(targetUid, principal.principalId(), "NOT_BOUND");
|
||||
break;
|
||||
}
|
||||
if (bindingService.isBlacklisted(targetUid, principal.principalId())) {
|
||||
sendError(session, "该设备已拒绝来自你的连接");
|
||||
metrics.incBlockedOffer();
|
||||
auditOfferBlocked(targetUid, principal.principalId(), "BLACKLISTED");
|
||||
break;
|
||||
}
|
||||
}
|
||||
// 连接请求:统一经 ConnectionRequestManager 做校验/去重/待确认跟踪后再转发
|
||||
handleConnectionRequest(signalMessage);
|
||||
break;
|
||||
@@ -110,13 +240,13 @@ public class SignalWebSocketHandler extends TextWebSocketHandler {
|
||||
break;
|
||||
}
|
||||
} catch (Exception e) {
|
||||
logger.error("Error handling message: {}", e.getMessage(), e);
|
||||
logger.error("处理消息出错: {}", e.getMessage(), e);
|
||||
}
|
||||
}
|
||||
|
||||
@Override
|
||||
public void afterConnectionClosed(WebSocketSession session, CloseStatus status) throws Exception {
|
||||
logger.info("WebSocket connection closed: {} ({})", session.getId(), status);
|
||||
logger.info("连接已关闭: {} ({})", session.getId(), status);
|
||||
sessionManager.unregisterSession(session);
|
||||
metrics.recordSessionCount(sessionManager.getOnlineCount());
|
||||
}
|
||||
@@ -131,10 +261,10 @@ public class SignalWebSocketHandler extends TextWebSocketHandler {
|
||||
|| exception.getMessage().contains("Broken pipe")
|
||||
|| exception.getMessage().contains("An established connection")));
|
||||
if (benign) {
|
||||
logger.debug("Transport closed (client disconnected) on session {}: {}",
|
||||
logger.debug("传输层关闭(客户端断开),会话 {}: {}",
|
||||
session.getId(), exception.getMessage());
|
||||
} else {
|
||||
logger.warn("Transport error on session {}: {}", session.getId(), exception.getMessage());
|
||||
logger.warn("传输层错误,会话 {}: {}", session.getId(), exception.getMessage());
|
||||
}
|
||||
sessionManager.unregisterSession(session);
|
||||
metrics.recordSessionCount(sessionManager.getOnlineCount());
|
||||
@@ -144,44 +274,64 @@ public class SignalWebSocketHandler extends TextWebSocketHandler {
|
||||
}
|
||||
}
|
||||
|
||||
private void handleRegister(WebSocketSession session, SignalMessage message) {
|
||||
String deviceId = message.getFromDeviceId();
|
||||
String deviceTypeStr = message.getDeviceType();
|
||||
|
||||
if (deviceId == null || deviceTypeStr == null) {
|
||||
logger.warn("Invalid REGISTER message: missing deviceId or deviceType");
|
||||
return;
|
||||
}
|
||||
|
||||
try {
|
||||
DeviceType deviceType = DeviceType.valueOf(deviceTypeStr.toUpperCase());
|
||||
sessionManager.registerDevice(deviceId, deviceType, session);
|
||||
metrics.incRegistered();
|
||||
metrics.recordSessionCount(sessionManager.getOnlineCount());
|
||||
|
||||
// 回复注册成功
|
||||
Map<String, Object> response = new HashMap<>();
|
||||
response.put("type", "REGISTER_SUCCESS");
|
||||
response.put("deviceId", deviceId);
|
||||
sendToSession(session, response);
|
||||
|
||||
logger.info("Device {} registered as {}", deviceId, deviceType);
|
||||
} catch (IllegalArgumentException e) {
|
||||
logger.warn("Invalid device type: {}", deviceTypeStr);
|
||||
}
|
||||
private AuthPrincipal principalOf(WebSocketSession session) {
|
||||
Object attr = session.getAttributes().get(AuthHandshakeInterceptor.ATTR_PRINCIPAL);
|
||||
return attr instanceof AuthPrincipal principal ? principal : null;
|
||||
}
|
||||
|
||||
private void handleDeviceList(WebSocketSession session, SignalMessage message) {
|
||||
List<String> controllers = sessionManager.getDevicesByType(DeviceType.CONTROLLER);
|
||||
List<String> controlled = sessionManager.getDevicesByType(DeviceType.CONTROLLED);
|
||||
/**
|
||||
* 兼容旧客户端的 REGISTER:不再接受客户端自报身份,仅回显服务端裁定结果。
|
||||
*/
|
||||
private void handleLegacyRegister(WebSocketSession session, AuthPrincipal principal) {
|
||||
Map<String, Object> response = new HashMap<>();
|
||||
response.put("type", "REGISTER_SUCCESS");
|
||||
response.put("deviceId", principal.deviceId());
|
||||
response.put("deviceType", principal.deviceType().name());
|
||||
response.put("displayName", principal.displayName());
|
||||
sendToSession(session, response);
|
||||
}
|
||||
|
||||
/**
|
||||
* 设备列表(绑定视图)。
|
||||
*
|
||||
* <p>为防止枚举被控端,服务端不再向任意主控端返回全局在线设备清单。
|
||||
* 主控端仅能看到「自己已绑定」的设备及其在线状态;被控端此项为空列表。
|
||||
*/
|
||||
private void handleDeviceList(WebSocketSession session, AuthPrincipal principal) {
|
||||
Map<String, Object> response = new HashMap<>();
|
||||
response.put("type", "DEVICE_LIST");
|
||||
response.put("controllers", controllers);
|
||||
response.put("controlled", controlled);
|
||||
|
||||
if (principal.deviceType() == DeviceType.CONTROLLER) {
|
||||
List<Map<String, Object>> devices = bindingService.listByUser(principal.principalId())
|
||||
.stream()
|
||||
.filter(b -> b.getStatus() == DeviceBinding.BindingStatus.ACTIVE)
|
||||
.map(b -> {
|
||||
Map<String, Object> item = new HashMap<>();
|
||||
item.put("deviceUid", b.getDeviceUid());
|
||||
item.put("alias", b.getAlias());
|
||||
item.put("role", b.getRole().name());
|
||||
item.put("online", sessionManager.isDeviceOnline(b.getDeviceUid()));
|
||||
return item;
|
||||
})
|
||||
.toList();
|
||||
response.put("boundDevices", devices);
|
||||
} else {
|
||||
// 被控端无需浏览主控端列表
|
||||
response.put("boundDevices", List.of());
|
||||
}
|
||||
sendToSession(session, response);
|
||||
}
|
||||
|
||||
/** 记录被拦截的 OFFER(未绑定/被拉黑),便于审计与风控。 */
|
||||
private void auditOfferBlocked(String deviceUid, String userId, String reason) {
|
||||
try {
|
||||
auditService.recordUser(userId, AuditService.ACTION_OFFER_BLOCKED,
|
||||
AuditService.RESULT_SUCCESS, null, "device=" + deviceUid + " reason=" + reason);
|
||||
} catch (Exception e) {
|
||||
logger.warn("审计 OFFER 拦截记录失败: {}", e.getMessage());
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* 处理主控端发来的连接请求(OFFER):先做设备类型校验与去重,
|
||||
* 再登记“待被控端确认”状态,最后转发给被控端。
|
||||
@@ -193,7 +343,7 @@ public class SignalWebSocketHandler extends TextWebSocketHandler {
|
||||
// 1. 校验:仅允许 CONTROLLER -> CONTROLLED
|
||||
String error = connectionRequestManager.validateOffer(fromDeviceId, toDeviceId);
|
||||
if (error != null) {
|
||||
logger.warn("Invalid connection request {} -> {}: {}", fromDeviceId, toDeviceId, error);
|
||||
logger.warn("非法连接请求 {} -> {}: {}", fromDeviceId, toDeviceId, error);
|
||||
Map<String, Object> response = new HashMap<>();
|
||||
response.put("type", "REQUEST_ERROR");
|
||||
response.put("toDeviceId", toDeviceId);
|
||||
@@ -225,13 +375,13 @@ public class SignalWebSocketHandler extends TextWebSocketHandler {
|
||||
|
||||
// 2. 去重:短时间内重复 OFFER 直接忽略,避免被控端反复弹窗
|
||||
if (connectionRequestManager.isDuplicateOffer(fromDeviceId, toDeviceId)) {
|
||||
logger.info("Duplicate connection request {} -> {} ignored", fromDeviceId, toDeviceId);
|
||||
logger.info("重复连接请求 {} -> {} 已忽略", fromDeviceId, toDeviceId);
|
||||
return;
|
||||
}
|
||||
|
||||
// 3. 目标不在线:立即回送 TARGET_OFFLINE,且不登记待确认(避免等待超时)
|
||||
if (!sessionManager.isDeviceOnline(toDeviceId)) {
|
||||
logger.warn("Target device {} is offline, cannot deliver connection request", toDeviceId);
|
||||
logger.warn("目标被控端 {} 不在线,无法投递连接请求", toDeviceId);
|
||||
notifySenderTargetOffline(message, toDeviceId);
|
||||
return;
|
||||
}
|
||||
@@ -244,13 +394,13 @@ public class SignalWebSocketHandler extends TextWebSocketHandler {
|
||||
private void forwardMessage(SignalMessage message, boolean notifyOffline) {
|
||||
String toDeviceId = message.getToDeviceId();
|
||||
if (toDeviceId == null) {
|
||||
logger.warn("Cannot forward message: toDeviceId is null");
|
||||
logger.warn("无法转发消息:toDeviceId 为空");
|
||||
return;
|
||||
}
|
||||
|
||||
WebSocketSession targetSession = sessionManager.getSession(toDeviceId);
|
||||
if (targetSession == null || !targetSession.isOpen()) {
|
||||
logger.warn("Target device {} is not online", toDeviceId);
|
||||
logger.warn("目标设备 {} 不在线", toDeviceId);
|
||||
// 仅 OFFER / ANSWER 在目标不在线时回送 TARGET_OFFLINE,其余类型静默丢弃
|
||||
if (notifyOffline) {
|
||||
notifySenderTargetOffline(message, toDeviceId);
|
||||
@@ -261,9 +411,9 @@ public class SignalWebSocketHandler extends TextWebSocketHandler {
|
||||
try {
|
||||
String jsonMessage = objectMapper.writeValueAsString(message);
|
||||
targetSession.sendMessage(new TextMessage(jsonMessage));
|
||||
logger.debug("Forwarded {} from {} to {}", message.getType(), message.getFromDeviceId(), toDeviceId);
|
||||
logger.debug("已转发 {}:{} -> {}", message.getType(), message.getFromDeviceId(), toDeviceId);
|
||||
} catch (IOException e) {
|
||||
logger.error("Error forwarding message to {}: {}", toDeviceId, e.getMessage());
|
||||
logger.error("转发消息到 {} 失败: {}", toDeviceId, e.getMessage());
|
||||
}
|
||||
}
|
||||
|
||||
@@ -273,23 +423,23 @@ public class SignalWebSocketHandler extends TextWebSocketHandler {
|
||||
private void notifySenderTargetOffline(SignalMessage message, String offlineDeviceId) {
|
||||
String fromDeviceId = message.getFromDeviceId();
|
||||
if (fromDeviceId == null) {
|
||||
logger.warn("Cannot notify offline state: fromDeviceId is null");
|
||||
logger.warn("无法通知离线状态:fromDeviceId 为空");
|
||||
return;
|
||||
}
|
||||
|
||||
WebSocketSession senderSession = sessionManager.getSession(fromDeviceId);
|
||||
if (senderSession == null || !senderSession.isOpen()) {
|
||||
logger.warn("Sender {} session not found, cannot notify target offline", fromDeviceId);
|
||||
logger.warn("发送方 {} 会话不存在,无法通知目标离线", fromDeviceId);
|
||||
return;
|
||||
}
|
||||
|
||||
Map<String, Object> response = new HashMap<>();
|
||||
response.put("type", "TARGET_OFFLINE");
|
||||
response.put("toDeviceId", offlineDeviceId);
|
||||
response.put("payload", "目标被控端(" + offlineDeviceId + ")不在线,请确认设备已开启并连接到信令服务器");
|
||||
response.put("payload", "目标被控端不在线,请确认设备已开启并连接到信令服务器");
|
||||
|
||||
sendToSession(senderSession, response);
|
||||
logger.info("Notified sender {} that target {} is offline", fromDeviceId, offlineDeviceId);
|
||||
logger.info("已通知发送方 {} 目标 {} 不在线", fromDeviceId, offlineDeviceId);
|
||||
}
|
||||
|
||||
private void sendToSession(WebSocketSession session, Object data) {
|
||||
@@ -297,7 +447,15 @@ public class SignalWebSocketHandler extends TextWebSocketHandler {
|
||||
String json = objectMapper.writeValueAsString(data);
|
||||
session.sendMessage(new TextMessage(json));
|
||||
} catch (IOException e) {
|
||||
logger.error("Error sending message to session {}: {}", session.getId(), e.getMessage());
|
||||
logger.error("向会话 {} 发送消息失败: {}", session.getId(), e.getMessage());
|
||||
}
|
||||
}
|
||||
|
||||
/** 向发送方回送一条 REQUEST_ERROR 业务提示(不关闭连接)。 */
|
||||
private void sendError(WebSocketSession session, String message) {
|
||||
Map<String, Object> response = new HashMap<>();
|
||||
response.put("type", "REQUEST_ERROR");
|
||||
response.put("payload", message);
|
||||
sendToSession(session, response);
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user