feat: 增加请求弹窗和拒绝说明
This commit is contained in:
@@ -1,6 +1,7 @@
|
||||
package com.ttstd.signaling.handler;
|
||||
|
||||
import com.fasterxml.jackson.databind.ObjectMapper;
|
||||
import com.ttstd.signaling.manager.ConnectionRequestManager;
|
||||
import com.ttstd.signaling.manager.SessionManager;
|
||||
import com.ttstd.signaling.model.DeviceType;
|
||||
import com.ttstd.signaling.model.SignalMessage;
|
||||
@@ -24,9 +25,24 @@ public class SignalWebSocketHandler extends TextWebSocketHandler {
|
||||
private static final ObjectMapper objectMapper = new ObjectMapper();
|
||||
|
||||
private final SessionManager sessionManager;
|
||||
private final ConnectionRequestManager connectionRequestManager;
|
||||
|
||||
public SignalWebSocketHandler(SessionManager sessionManager) {
|
||||
public SignalWebSocketHandler(SessionManager sessionManager, ConnectionRequestManager connectionRequestManager) {
|
||||
this.sessionManager = sessionManager;
|
||||
this.connectionRequestManager = connectionRequestManager;
|
||||
this.connectionRequestManager.setSender(this::sendToDevice);
|
||||
}
|
||||
|
||||
/**
|
||||
* 向指定设备回送服务端事件(错误/超时通知等)。
|
||||
*/
|
||||
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);
|
||||
return;
|
||||
}
|
||||
sendToSession(session, message);
|
||||
}
|
||||
|
||||
@Override
|
||||
@@ -56,10 +72,21 @@ public class SignalWebSocketHandler extends TextWebSocketHandler {
|
||||
handleDeviceList(session, signalMessage);
|
||||
break;
|
||||
case "OFFER":
|
||||
// 连接请求:统一经 ConnectionRequestManager 做校验/去重/待确认跟踪后再转发
|
||||
handleConnectionRequest(signalMessage);
|
||||
break;
|
||||
case "ANSWER":
|
||||
// 关键信令:目标不在线时回送 TARGET_OFFLINE 提醒发送方
|
||||
// 被控端已接受:清理待确认状态并转发 Answer 给主控端
|
||||
connectionRequestManager.completePendingRequest(
|
||||
signalMessage.getToDeviceId(), signalMessage.getFromDeviceId());
|
||||
forwardMessage(signalMessage, true);
|
||||
break;
|
||||
case "CONNECTION_REJECTED":
|
||||
// 被控端已拒绝:清理待确认状态并转发拒绝通知给主控端
|
||||
connectionRequestManager.completePendingRequest(
|
||||
signalMessage.getToDeviceId(), signalMessage.getFromDeviceId());
|
||||
forwardMessage(signalMessage, false);
|
||||
break;
|
||||
case "ICE_CANDIDATE":
|
||||
case "CONTROL_COMMAND":
|
||||
// 普通信令:目标不在线时静默丢弃,不返回提示
|
||||
@@ -126,6 +153,44 @@ public class SignalWebSocketHandler extends TextWebSocketHandler {
|
||||
sendToSession(session, response);
|
||||
}
|
||||
|
||||
/**
|
||||
* 处理主控端发来的连接请求(OFFER):先做设备类型校验与去重,
|
||||
* 再登记“待被控端确认”状态,最后转发给被控端。
|
||||
*/
|
||||
private void handleConnectionRequest(SignalMessage message) {
|
||||
String fromDeviceId = message.getFromDeviceId();
|
||||
String toDeviceId = message.getToDeviceId();
|
||||
|
||||
// 1. 校验:仅允许 CONTROLLER -> CONTROLLED
|
||||
String error = connectionRequestManager.validateOffer(fromDeviceId, toDeviceId);
|
||||
if (error != null) {
|
||||
logger.warn("Invalid connection request {} -> {}: {}", fromDeviceId, toDeviceId, error);
|
||||
Map<String, Object> response = new HashMap<>();
|
||||
response.put("type", "REQUEST_ERROR");
|
||||
response.put("toDeviceId", toDeviceId);
|
||||
response.put("payload", error);
|
||||
sendToDevice(fromDeviceId, response);
|
||||
return;
|
||||
}
|
||||
|
||||
// 2. 去重:短时间内重复 OFFER 直接忽略,避免被控端反复弹窗
|
||||
if (connectionRequestManager.isDuplicateOffer(fromDeviceId, toDeviceId)) {
|
||||
logger.info("Duplicate connection request {} -> {} ignored", fromDeviceId, toDeviceId);
|
||||
return;
|
||||
}
|
||||
|
||||
// 3. 目标不在线:立即回送 TARGET_OFFLINE,且不登记待确认(避免等待超时)
|
||||
if (!sessionManager.isDeviceOnline(toDeviceId)) {
|
||||
logger.warn("Target device {} is offline, cannot deliver connection request", toDeviceId);
|
||||
notifySenderTargetOffline(message, toDeviceId);
|
||||
return;
|
||||
}
|
||||
|
||||
// 4. 登记待确认并转发 OFFER 给被控端
|
||||
connectionRequestManager.registerPendingRequest(fromDeviceId, toDeviceId);
|
||||
forwardMessage(message, false);
|
||||
}
|
||||
|
||||
private void forwardMessage(SignalMessage message, boolean notifyOffline) {
|
||||
String toDeviceId = message.getToDeviceId();
|
||||
if (toDeviceId == null) {
|
||||
|
||||
@@ -0,0 +1,153 @@
|
||||
package com.ttstd.signaling.manager;
|
||||
|
||||
import com.ttstd.signaling.model.DeviceType;
|
||||
import org.slf4j.Logger;
|
||||
import org.slf4j.LoggerFactory;
|
||||
import org.springframework.stereotype.Component;
|
||||
|
||||
import java.util.Map;
|
||||
import java.util.concurrent.ConcurrentHashMap;
|
||||
import java.util.concurrent.Executors;
|
||||
import java.util.concurrent.ScheduledExecutorService;
|
||||
import java.util.concurrent.TimeUnit;
|
||||
|
||||
/**
|
||||
* 连接请求管理器:对主控端发往被控端的 OFFER(连接请求)进行统一治理。
|
||||
*
|
||||
* <p>职责:
|
||||
* <ol>
|
||||
* <li><b>校验</b>:OFFER 仅允许由 CONTROLLER 发往 CONTROLLED,否则拒绝并向发送方报错。</li>
|
||||
* <li><b>去重</b>:同一主控端对同一被控端的重复 OFFER(如网络重传/重复点击)在短时间内被忽略,
|
||||
* 避免被控端反复弹窗。</li>
|
||||
* <li><b>待确认跟踪</b>:记录处于“等待被控端确认”状态的请求,收到 ANSWER(接受)或
|
||||
* CONNECTION_REJECTED(拒绝)时清理;超时未决则通知主控端请求超时。</li>
|
||||
* </ol>
|
||||
*/
|
||||
@Component
|
||||
public class ConnectionRequestManager {
|
||||
|
||||
private static final Logger logger = LoggerFactory.getLogger(ConnectionRequestManager.class);
|
||||
|
||||
/** 同一连接的去重窗口:窗口内重复的 OFFER 将被忽略。 */
|
||||
private static final long DEDUP_WINDOW_MS = 5_000;
|
||||
/** 被控端确认超时:超过该时间既无 ANSWER 也无拒绝,则通知主控端超时。 */
|
||||
private static final long REQUEST_TIMEOUT_MS = 60_000;
|
||||
|
||||
private final SessionManager sessionManager;
|
||||
private final ScheduledExecutorService scheduler = Executors.newSingleThreadScheduledExecutor(r -> {
|
||||
Thread t = new Thread(r, "connection-request-timeout");
|
||||
t.setDaemon(true);
|
||||
return t;
|
||||
});
|
||||
|
||||
/** 最近一次 OFFER 的时间戳,键为 "fromDeviceId->toDeviceId"。 */
|
||||
private final Map<String, Long> lastOfferTime = new ConcurrentHashMap<>();
|
||||
/** 处于“等待被控端确认”状态的请求,键同上。 */
|
||||
private final Map<String, PendingRequest> pendingRequests = new ConcurrentHashMap<>();
|
||||
|
||||
public ConnectionRequestManager(SessionManager sessionManager) {
|
||||
this.sessionManager = sessionManager;
|
||||
}
|
||||
|
||||
/** 向指定设备发送消息的回调(由 Handler 注入,用于回送错误/超时通知)。 */
|
||||
public interface SignalSender {
|
||||
void send(String deviceId, Object message);
|
||||
}
|
||||
|
||||
private SignalSender sender;
|
||||
|
||||
public void setSender(SignalSender sender) {
|
||||
this.sender = sender;
|
||||
}
|
||||
|
||||
private static String keyOf(String from, String to) {
|
||||
return from + "->" + to;
|
||||
}
|
||||
|
||||
/**
|
||||
* 校验 OFFER 的合法性(发送方必须为 CONTROLLER,接收方必须为 CONTROLLED)。
|
||||
*
|
||||
* @return 校验失败时的错误信息(含发送目标),成功返回 null。
|
||||
*/
|
||||
public String validateOffer(String fromDeviceId, String toDeviceId) {
|
||||
if (fromDeviceId == null || toDeviceId == null) {
|
||||
return "连接请求缺少发送方或接收方设备ID";
|
||||
}
|
||||
DeviceType fromType = sessionManager.getDeviceType(fromDeviceId);
|
||||
DeviceType toType = sessionManager.getDeviceType(toDeviceId);
|
||||
if (fromType != DeviceType.CONTROLLER) {
|
||||
return "连接请求只能由主控端(CONTROLLER)发起";
|
||||
}
|
||||
if (toType != DeviceType.CONTROLLED) {
|
||||
return "连接请求的目标必须是被控端(CONTROLLED)";
|
||||
}
|
||||
return null;
|
||||
}
|
||||
|
||||
/**
|
||||
* 去重判断:若同一连接的最近一次 OFFER 在去重窗口内,则视为重复。
|
||||
*/
|
||||
public boolean isDuplicateOffer(String fromDeviceId, String toDeviceId) {
|
||||
String key = keyOf(fromDeviceId, toDeviceId);
|
||||
long now = System.currentTimeMillis();
|
||||
Long last = lastOfferTime.get(key);
|
||||
if (last != null && (now - last) < DEDUP_WINDOW_MS) {
|
||||
return true;
|
||||
}
|
||||
lastOfferTime.put(key, now);
|
||||
return false;
|
||||
}
|
||||
|
||||
/**
|
||||
* 登记一个待确认的连接请求,并安排超时任务。
|
||||
*/
|
||||
public void registerPendingRequest(String fromDeviceId, String toDeviceId) {
|
||||
String key = keyOf(fromDeviceId, toDeviceId);
|
||||
PendingRequest pending = new PendingRequest(fromDeviceId, toDeviceId);
|
||||
pending.timeoutTask = scheduler.schedule(() -> {
|
||||
// 超时:仅当该请求仍未被接受/拒绝时通知主控端。
|
||||
if (pendingRequests.remove(key, pending)) {
|
||||
logger.warn("Connection request {} timed out after {}ms", key, REQUEST_TIMEOUT_MS);
|
||||
if (sender != null) {
|
||||
sender.send(fromDeviceId, buildEvent("REQUEST_TIMEOUT", toDeviceId,
|
||||
"被控端(" + toDeviceId + ")未在限定时间内响应连接请求,请稍后重试"));
|
||||
}
|
||||
}
|
||||
}, REQUEST_TIMEOUT_MS, TimeUnit.MILLISECONDS);
|
||||
pendingRequests.put(key, pending);
|
||||
logger.info("Registered pending connection request {}", key);
|
||||
}
|
||||
|
||||
/**
|
||||
* 当被控端接受(收到 ANSWER)或拒绝(收到 CONNECTION_REJECTED)时,清理待确认状态。
|
||||
*/
|
||||
public void completePendingRequest(String fromDeviceId, String toDeviceId) {
|
||||
String key = keyOf(fromDeviceId, toDeviceId);
|
||||
PendingRequest pending = pendingRequests.remove(key);
|
||||
if (pending != null) {
|
||||
if (pending.timeoutTask != null) {
|
||||
pending.timeoutTask.cancel(false);
|
||||
}
|
||||
logger.info("Completed pending connection request {}", key);
|
||||
}
|
||||
}
|
||||
|
||||
private Map<String, Object> buildEvent(String type, String toDeviceId, String message) {
|
||||
Map<String, Object> event = new ConcurrentHashMap<>();
|
||||
event.put("type", type);
|
||||
event.put("toDeviceId", toDeviceId);
|
||||
event.put("payload", message);
|
||||
return event;
|
||||
}
|
||||
|
||||
private static class PendingRequest {
|
||||
final String fromDeviceId;
|
||||
final String toDeviceId;
|
||||
java.util.concurrent.ScheduledFuture<?> timeoutTask;
|
||||
|
||||
PendingRequest(String fromDeviceId, String toDeviceId) {
|
||||
this.fromDeviceId = fromDeviceId;
|
||||
this.toDeviceId = toDeviceId;
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -48,6 +48,10 @@ public class SessionManager {
|
||||
return session != null && session.isOpen();
|
||||
}
|
||||
|
||||
public DeviceType getDeviceType(String deviceId) {
|
||||
return deviceTypes.get(deviceId);
|
||||
}
|
||||
|
||||
public List<String> getDevicesByType(DeviceType type) {
|
||||
List<String> result = new ArrayList<>();
|
||||
for (Map.Entry<String, DeviceType> entry : deviceTypes.entrySet()) {
|
||||
|
||||
Reference in New Issue
Block a user