|
|
|
@ -41,16 +41,25 @@ import org.thingsboard.server.service.telemetry.TelemetryWebSocketMsgEndpoint; |
|
|
|
import org.thingsboard.server.service.telemetry.TelemetryWebSocketService; |
|
|
|
import org.thingsboard.server.service.telemetry.TelemetryWebSocketSessionRef; |
|
|
|
|
|
|
|
import javax.websocket.*; |
|
|
|
import javax.annotation.PreDestroy; |
|
|
|
import javax.websocket.RemoteEndpoint; |
|
|
|
import javax.websocket.SendHandler; |
|
|
|
import javax.websocket.SendResult; |
|
|
|
import javax.websocket.Session; |
|
|
|
import java.io.IOException; |
|
|
|
import java.net.URI; |
|
|
|
import java.nio.ByteBuffer; |
|
|
|
import java.security.InvalidParameterException; |
|
|
|
import java.util.Queue; |
|
|
|
import java.util.Set; |
|
|
|
import java.util.UUID; |
|
|
|
import java.util.concurrent.ConcurrentHashMap; |
|
|
|
import java.util.concurrent.ConcurrentMap; |
|
|
|
import java.util.concurrent.Executors; |
|
|
|
import java.util.concurrent.LinkedBlockingQueue; |
|
|
|
import java.util.concurrent.ScheduledExecutorService; |
|
|
|
import java.util.concurrent.TimeUnit; |
|
|
|
import java.util.concurrent.atomic.AtomicLong; |
|
|
|
|
|
|
|
@Service |
|
|
|
@TbCoreComponent |
|
|
|
@ -79,6 +88,9 @@ public class TbWebSocketHandler extends TextWebSocketHandler implements Telemetr |
|
|
|
@Value("${server.ws.limits.max_updates_per_session:}") |
|
|
|
private String perSessionUpdatesConfiguration; |
|
|
|
|
|
|
|
@Value("${server.ws.ping_timeout:30000}") |
|
|
|
private long pingTimeout; |
|
|
|
|
|
|
|
private ConcurrentMap<String, TelemetryWebSocketSessionRef> blacklistedSessions = new ConcurrentHashMap<>(); |
|
|
|
private ConcurrentMap<String, TbRateLimits> perSessionUpdateLimits = new ConcurrentHashMap<>(); |
|
|
|
|
|
|
|
@ -87,6 +99,8 @@ public class TbWebSocketHandler extends TextWebSocketHandler implements Telemetr |
|
|
|
private ConcurrentMap<UserId, Set<String>> regularUserSessionsMap = new ConcurrentHashMap<>(); |
|
|
|
private ConcurrentMap<UserId, Set<String>> publicUserSessionsMap = new ConcurrentHashMap<>(); |
|
|
|
|
|
|
|
private ScheduledExecutorService pingExecutor = Executors.newSingleThreadScheduledExecutor(); |
|
|
|
|
|
|
|
@Override |
|
|
|
public void handleTextMessage(WebSocketSession session, TextMessage message) { |
|
|
|
try { |
|
|
|
@ -120,6 +134,9 @@ public class TbWebSocketHandler extends TextWebSocketHandler implements Telemetr |
|
|
|
return; |
|
|
|
} |
|
|
|
internalSessionMap.put(internalSessionId, new SessionMetaData(session, sessionRef, maxMsgQueuePerSession)); |
|
|
|
|
|
|
|
pingExecutor.scheduleWithFixedDelay(() -> internalSessionMap.get(internalSessionId).sendPing(), 10000, 10000, TimeUnit.MILLISECONDS); |
|
|
|
|
|
|
|
externalSessionMap.put(externalSessionId, internalSessionId); |
|
|
|
processInWebSocketService(sessionRef, SessionEvent.onEstablished()); |
|
|
|
log.info("[{}][{}][{}] Session is opened", sessionRef.getSecurityCtx().getTenantId(), externalSessionId, session.getId()); |
|
|
|
@ -189,6 +206,8 @@ public class TbWebSocketHandler extends TextWebSocketHandler implements Telemetr |
|
|
|
private volatile boolean isSending = false; |
|
|
|
private final Queue<String> msgQueue; |
|
|
|
|
|
|
|
private final AtomicLong lastActivityTime; |
|
|
|
|
|
|
|
SessionMetaData(WebSocketSession session, TelemetryWebSocketSessionRef sessionRef, int maxMsgQueuePerSession) { |
|
|
|
super(); |
|
|
|
this.session = session; |
|
|
|
@ -196,12 +215,30 @@ public class TbWebSocketHandler extends TextWebSocketHandler implements Telemetr |
|
|
|
this.asyncRemote = nativeSession.getAsyncRemote(); |
|
|
|
this.sessionRef = sessionRef; |
|
|
|
this.msgQueue = new LinkedBlockingQueue<>(maxMsgQueuePerSession); |
|
|
|
this.lastActivityTime = new AtomicLong(System.currentTimeMillis()); |
|
|
|
} |
|
|
|
|
|
|
|
synchronized void sendPing() { |
|
|
|
try { |
|
|
|
if (System.currentTimeMillis() - lastActivityTime.get() >= pingTimeout) { |
|
|
|
this.asyncRemote.sendPing(ByteBuffer.wrap(new byte[]{})); |
|
|
|
lastActivityTime.set(System.currentTimeMillis()); |
|
|
|
} |
|
|
|
} catch (Exception e) { |
|
|
|
log.trace("[{}] Failed to send ping msg", session.getId(), e); |
|
|
|
try { |
|
|
|
close(this.sessionRef, CloseStatus.SESSION_NOT_RELIABLE); |
|
|
|
} catch (IOException ioe) { |
|
|
|
log.trace("[{}] Session transport error", session.getId(), ioe); |
|
|
|
} |
|
|
|
} |
|
|
|
} |
|
|
|
|
|
|
|
synchronized void sendMsg(String msg) { |
|
|
|
if (isSending) { |
|
|
|
try { |
|
|
|
msgQueue.add(msg); |
|
|
|
lastActivityTime.set(System.currentTimeMillis()); |
|
|
|
} catch (RuntimeException e) { |
|
|
|
if (log.isTraceEnabled()) { |
|
|
|
log.trace("[{}][{}] Session closed due to queue error", sessionRef.getSecurityCtx().getTenantId(), session.getId(), e); |
|
|
|
@ -393,4 +430,11 @@ public class TbWebSocketHandler extends TextWebSocketHandler implements Telemetr |
|
|
|
} |
|
|
|
} |
|
|
|
|
|
|
|
@PreDestroy |
|
|
|
private void destroy() { |
|
|
|
if (pingExecutor != null) { |
|
|
|
pingExecutor.shutdownNow(); |
|
|
|
} |
|
|
|
} |
|
|
|
|
|
|
|
} |