Browse Source

Added functionality to handle pong responses on web sockets

pull/5793/head
Volodymyr Babak 5 years ago
parent
commit
7570ab5703
  1. 57
      application/src/main/java/org/thingsboard/server/controller/plugin/TbWebSocketHandler.java
  2. 5
      application/src/main/java/org/thingsboard/server/service/telemetry/DefaultTelemetryWebSocketService.java

57
application/src/main/java/org/thingsboard/server/controller/plugin/TbWebSocketHandler.java

@ -23,6 +23,7 @@ import org.springframework.security.core.Authentication;
import org.springframework.stereotype.Service;
import org.springframework.util.StringUtils;
import org.springframework.web.socket.CloseStatus;
import org.springframework.web.socket.PongMessage;
import org.springframework.web.socket.TextMessage;
import org.springframework.web.socket.WebSocketSession;
import org.springframework.web.socket.adapter.NativeWebSocketSession;
@ -111,6 +112,22 @@ public class TbWebSocketHandler extends TextWebSocketHandler implements Telemetr
}
}
@Override
protected void handlePongMessage(WebSocketSession session, PongMessage message) throws Exception {
try {
SessionMetaData sessionMd = internalSessionMap.get(session.getId());
if (sessionMd != null) {
log.trace("[{}][{}] Processing pong response {}", sessionMd.sessionRef.getSecurityCtx().getTenantId(), session.getId(), message.getPayload());
sessionMd.processPongMessage(System.currentTimeMillis());
} else {
log.trace("[{}] Failed to find session", session.getId());
session.close(CloseStatus.SERVER_ERROR.withReason("Session not found!"));
}
} catch (IOException e) {
log.warn("IO error", e);
}
}
@Override
public void afterConnectionEstablished(WebSocketSession session) throws Exception {
super.afterConnectionEstablished(session);
@ -213,19 +230,29 @@ public class TbWebSocketHandler extends TextWebSocketHandler implements Telemetr
synchronized void sendPing(long currentTime) {
try {
if (currentTime - lastActivityTime >= pingTimeout) {
log.warn("[{}] Closing session due to ping timeout", session.getId());
closeSession(CloseStatus.SESSION_NOT_RELIABLE);
} else {
this.asyncRemote.sendPing(PING_MSG);
lastActivityTime = currentTime;
}
} 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);
}
closeSession(CloseStatus.SESSION_NOT_RELIABLE);
}
}
private void closeSession(CloseStatus reason) {
try {
close(this.sessionRef, reason);
} catch (IOException ioe) {
log.trace("[{}] Session transport error", session.getId(), ioe);
}
}
synchronized void processPongMessage(long currentTime) {
lastActivityTime = currentTime;
}
synchronized void sendMsg(String msg) {
if (isSending) {
try {
@ -236,11 +263,7 @@ public class TbWebSocketHandler extends TextWebSocketHandler implements Telemetr
} else {
log.info("[{}][{}] Session closed due to queue error", sessionRef.getSecurityCtx().getTenantId(), session.getId());
}
try {
close(sessionRef, CloseStatus.POLICY_VIOLATION.withReason("Max pending updates limit reached!"));
} catch (IOException ioe) {
log.trace("[{}] Session transport error", session.getId(), ioe);
}
closeSession(CloseStatus.POLICY_VIOLATION.withReason("Max pending updates limit reached!"));
}
} else {
isSending = true;
@ -253,11 +276,7 @@ public class TbWebSocketHandler extends TextWebSocketHandler implements Telemetr
this.asyncRemote.sendText(msg, this);
} catch (Exception e) {
log.trace("[{}] Failed to send msg", session.getId(), e);
try {
close(this.sessionRef, CloseStatus.SESSION_NOT_RELIABLE);
} catch (IOException ioe) {
log.trace("[{}] Session transport error", session.getId(), ioe);
}
closeSession(CloseStatus.SESSION_NOT_RELIABLE);
}
}
@ -265,11 +284,7 @@ public class TbWebSocketHandler extends TextWebSocketHandler implements Telemetr
public void onResult(SendResult result) {
if (!result.isOK()) {
log.trace("[{}] Failed to send msg", session.getId(), result.getException());
try {
close(this.sessionRef, CloseStatus.SESSION_NOT_RELIABLE);
} catch (IOException ioe) {
log.trace("[{}] Session transport error", session.getId(), ioe);
}
closeSession(CloseStatus.SESSION_NOT_RELIABLE);
} else {
lastActivityTime = System.currentTimeMillis();
String msg = msgQueue.poll();

5
application/src/main/java/org/thingsboard/server/service/telemetry/DefaultTelemetryWebSocketService.java

@ -146,6 +146,9 @@ public class DefaultTelemetryWebSocketService implements TelemetryWebSocketServi
@Value("${server.ws.limits.max_subscriptions_per_public_user:0}")
private int maxSubscriptionsPerPublicUser;
@Value("${server.ws.ping_timeout:30000}")
private long pingTimeout;
private ConcurrentMap<TenantId, Set<String>> tenantSubscriptionsMap = new ConcurrentHashMap<>();
private ConcurrentMap<CustomerId, Set<String>> customerSubscriptionsMap = new ConcurrentHashMap<>();
private ConcurrentMap<UserId, Set<String>> regularUserSubscriptionsMap = new ConcurrentHashMap<>();
@ -162,7 +165,7 @@ public class DefaultTelemetryWebSocketService implements TelemetryWebSocketServi
executor = ThingsBoardExecutors.newWorkStealingPool(50, getClass());
pingExecutor = Executors.newSingleThreadScheduledExecutor(ThingsBoardThreadFactory.forName("telemetry-web-socket-ping"));
pingExecutor.scheduleWithFixedDelay(this::sendPing, 10000, 10000, TimeUnit.MILLISECONDS);
pingExecutor.scheduleWithFixedDelay(this::sendPing, pingTimeout/3, pingTimeout/3, TimeUnit.MILLISECONDS);
}
@PreDestroy

Loading…
Cancel
Save