|
|
|
@ -95,10 +95,12 @@ import org.thingsboard.server.queue.util.TbTransportComponent; |
|
|
|
import javax.annotation.PostConstruct; |
|
|
|
import javax.annotation.PreDestroy; |
|
|
|
import java.util.Collections; |
|
|
|
import java.util.HashSet; |
|
|
|
import java.util.List; |
|
|
|
import java.util.Map; |
|
|
|
import java.util.Optional; |
|
|
|
import java.util.Random; |
|
|
|
import java.util.Set; |
|
|
|
import java.util.UUID; |
|
|
|
import java.util.concurrent.ConcurrentHashMap; |
|
|
|
import java.util.concurrent.ConcurrentMap; |
|
|
|
@ -162,6 +164,7 @@ public class DefaultTransportService implements TransportService { |
|
|
|
private ExecutorService mainConsumerExecutor; |
|
|
|
|
|
|
|
private final ConcurrentMap<UUID, SessionMetaData> sessions = new ConcurrentHashMap<>(); |
|
|
|
private final ConcurrentMap<UUID, SessionActivityData> sessionsActivity = new ConcurrentHashMap<>(); |
|
|
|
private final Map<String, RpcRequestMetadata> toServerRpcPendingMap = new ConcurrentHashMap<>(); |
|
|
|
|
|
|
|
private volatile boolean stopped = false; |
|
|
|
@ -545,8 +548,11 @@ public class DefaultTransportService implements TransportService { |
|
|
|
@Override |
|
|
|
public void process(TransportProtos.SessionInfoProto sessionInfo, TransportProtos.SubscribeToAttributeUpdatesMsg msg, TransportServiceCallback<Void> callback) { |
|
|
|
if (checkLimits(sessionInfo, msg, callback)) { |
|
|
|
SessionMetaData sessionMetaData = reportActivityInternal(sessionInfo); |
|
|
|
sessionMetaData.setSubscribedToAttributes(!msg.getUnsubscribe()); |
|
|
|
SessionMetaData sessionMetaData = sessions.get(toSessionId(sessionInfo)); |
|
|
|
if (sessionMetaData != null) { |
|
|
|
sessionMetaData.setSubscribedToAttributes(!msg.getUnsubscribe()); |
|
|
|
} |
|
|
|
reportActivityInternal(sessionInfo); |
|
|
|
sendToDeviceActor(sessionInfo, TransportToDeviceActorMsg.newBuilder().setSessionInfo(sessionInfo).setSubscribeToAttributes(msg).build(), |
|
|
|
new ApiStatsProxyCallback<>(getTenantId(sessionInfo), getCustomerId(sessionInfo), 1, callback)); |
|
|
|
} |
|
|
|
@ -555,8 +561,11 @@ public class DefaultTransportService implements TransportService { |
|
|
|
@Override |
|
|
|
public void process(TransportProtos.SessionInfoProto sessionInfo, TransportProtos.SubscribeToRPCMsg msg, TransportServiceCallback<Void> callback) { |
|
|
|
if (checkLimits(sessionInfo, msg, callback)) { |
|
|
|
SessionMetaData sessionMetaData = reportActivityInternal(sessionInfo); |
|
|
|
sessionMetaData.setSubscribedToRPC(!msg.getUnsubscribe()); |
|
|
|
SessionMetaData sessionMetaData = sessions.get(toSessionId(sessionInfo)); |
|
|
|
if (sessionMetaData != null) { |
|
|
|
sessionMetaData.setSubscribedToRPC(!msg.getUnsubscribe()); |
|
|
|
} |
|
|
|
reportActivityInternal(sessionInfo); |
|
|
|
sendToDeviceActor(sessionInfo, TransportToDeviceActorMsg.newBuilder().setSessionInfo(sessionInfo).setSubscribeToRPC(msg).build(), |
|
|
|
new ApiStatsProxyCallback<>(getTenantId(sessionInfo), getCustomerId(sessionInfo), 1, callback)); |
|
|
|
} |
|
|
|
@ -668,52 +677,61 @@ public class DefaultTransportService implements TransportService { |
|
|
|
} |
|
|
|
|
|
|
|
@Override |
|
|
|
public SessionMetaData reportActivity(TransportProtos.SessionInfoProto sessionInfo) { |
|
|
|
return reportActivityInternal(sessionInfo); |
|
|
|
public void reportActivity(TransportProtos.SessionInfoProto sessionInfo) { |
|
|
|
reportActivityInternal(sessionInfo); |
|
|
|
} |
|
|
|
|
|
|
|
private SessionMetaData reportActivityInternal(TransportProtos.SessionInfoProto sessionInfo) { |
|
|
|
private void reportActivityInternal(TransportProtos.SessionInfoProto sessionInfo) { |
|
|
|
UUID sessionId = toSessionId(sessionInfo); |
|
|
|
SessionMetaData sessionMetaData = sessions.get(sessionId); |
|
|
|
if (sessionMetaData != null) { |
|
|
|
sessionMetaData.updateLastActivityTime(); |
|
|
|
} |
|
|
|
return sessionMetaData; |
|
|
|
SessionActivityData sessionMetaData = sessionsActivity.computeIfAbsent(sessionId, id -> new SessionActivityData(sessionInfo)); |
|
|
|
sessionMetaData.updateLastActivityTime(); |
|
|
|
} |
|
|
|
|
|
|
|
private void checkInactivityAndReportActivity() { |
|
|
|
long expTime = System.currentTimeMillis() - sessionInactivityTimeout; |
|
|
|
sessions.forEach((uuid, sessionMD) -> { |
|
|
|
long lastActivityTime = sessionMD.getLastActivityTime(); |
|
|
|
TransportProtos.SessionInfoProto sessionInfo = sessionMD.getSessionInfo(); |
|
|
|
if (sessionInfo.getGwSessionIdMSB() != 0 && |
|
|
|
sessionInfo.getGwSessionIdLSB() != 0) { |
|
|
|
SessionMetaData gwMetaData = sessions.get(new UUID(sessionInfo.getGwSessionIdMSB(), sessionInfo.getGwSessionIdLSB())); |
|
|
|
Set<UUID> sessionsToRemove = new HashSet<>(); |
|
|
|
sessionsActivity.forEach((uuid, sessionAD) -> { |
|
|
|
long lastActivityTime = sessionAD.getLastActivityTime(); |
|
|
|
SessionMetaData sessionMD = sessions.get(uuid); |
|
|
|
if (sessionMD != null) { |
|
|
|
sessionAD.setSessionInfo(sessionMD.getSessionInfo()); |
|
|
|
} else { |
|
|
|
sessionsToRemove.add(uuid); |
|
|
|
} |
|
|
|
TransportProtos.SessionInfoProto sessionInfo = sessionAD.getSessionInfo(); |
|
|
|
|
|
|
|
if (sessionInfo.getGwSessionIdMSB() != 0 && sessionInfo.getGwSessionIdLSB() != 0) { |
|
|
|
var gwSessionId = new UUID(sessionInfo.getGwSessionIdMSB(), sessionInfo.getGwSessionIdLSB()); |
|
|
|
SessionMetaData gwMetaData = sessions.get(gwSessionId); |
|
|
|
SessionActivityData gwActivityData = sessionsActivity.get(gwSessionId); |
|
|
|
if (gwMetaData != null && gwMetaData.isOverwriteActivityTime()) { |
|
|
|
lastActivityTime = Math.max(gwMetaData.getLastActivityTime(), lastActivityTime); |
|
|
|
lastActivityTime = Math.max(gwActivityData.getLastActivityTime(), lastActivityTime); |
|
|
|
} |
|
|
|
} |
|
|
|
if (lastActivityTime < expTime) { |
|
|
|
if (log.isDebugEnabled()) { |
|
|
|
log.debug("[{}] Session has expired due to last activity time: {}", toSessionId(sessionInfo), lastActivityTime); |
|
|
|
if (sessionMD != null) { |
|
|
|
if (log.isDebugEnabled()) { |
|
|
|
log.debug("[{}] Session has expired due to last activity time: {}", toSessionId(sessionInfo), lastActivityTime); |
|
|
|
} |
|
|
|
sessions.remove(uuid); |
|
|
|
sessionsToRemove.add(uuid); |
|
|
|
process(sessionInfo, getSessionEventMsg(TransportProtos.SessionEvent.CLOSED), null); |
|
|
|
TransportProtos.SessionCloseNotificationProto sessionCloseNotificationProto = TransportProtos.SessionCloseNotificationProto |
|
|
|
.newBuilder() |
|
|
|
.setMessage("Session has expired due to last activity time!") |
|
|
|
.build(); |
|
|
|
sessionMD.getListener().onRemoteSessionCloseCommand(uuid, sessionCloseNotificationProto); |
|
|
|
} |
|
|
|
process(sessionInfo, getSessionEventMsg(TransportProtos.SessionEvent.CLOSED), null); |
|
|
|
sessions.remove(uuid); |
|
|
|
TransportProtos.SessionCloseNotificationProto sessionCloseNotificationProto = TransportProtos.SessionCloseNotificationProto |
|
|
|
.newBuilder() |
|
|
|
.setMessage("session has expired due to last activity time!") |
|
|
|
.build(); |
|
|
|
sessionMD.getListener().onRemoteSessionCloseCommand(uuid, sessionCloseNotificationProto); |
|
|
|
} else { |
|
|
|
if (lastActivityTime > sessionMD.getLastReportedActivityTime()) { |
|
|
|
if (lastActivityTime > sessionAD.getLastReportedActivityTime()) { |
|
|
|
final long lastActivityTimeFinal = lastActivityTime; |
|
|
|
process(sessionInfo, TransportProtos.SubscriptionInfoProto.newBuilder() |
|
|
|
.setAttributeSubscription(sessionMD.isSubscribedToAttributes()) |
|
|
|
.setRpcSubscription(sessionMD.isSubscribedToRPC()) |
|
|
|
.setAttributeSubscription(sessionMD != null && sessionMD.isSubscribedToAttributes()) |
|
|
|
.setRpcSubscription(sessionMD != null && sessionMD.isSubscribedToRPC()) |
|
|
|
.setLastActivityTime(lastActivityTime).build(), new TransportServiceCallback<Void>() { |
|
|
|
@Override |
|
|
|
public void onSuccess(Void msg) { |
|
|
|
sessionMD.setLastReportedActivityTime(lastActivityTimeFinal); |
|
|
|
sessionAD.setLastReportedActivityTime(lastActivityTimeFinal); |
|
|
|
} |
|
|
|
|
|
|
|
@Override |
|
|
|
@ -724,6 +742,8 @@ public class DefaultTransportService implements TransportService { |
|
|
|
} |
|
|
|
} |
|
|
|
}); |
|
|
|
// Removes all closed or short-lived sessions.
|
|
|
|
sessionsToRemove.forEach(sessionsActivity::remove); |
|
|
|
} |
|
|
|
|
|
|
|
@Override |
|
|
|
@ -1146,4 +1166,9 @@ public class DefaultTransportService implements TransportService { |
|
|
|
public ExecutorService getCallbackExecutor() { |
|
|
|
return transportCallbackExecutor; |
|
|
|
} |
|
|
|
|
|
|
|
@Override |
|
|
|
public boolean hasSession(TransportProtos.SessionInfoProto sessionInfo) { |
|
|
|
return sessions.containsKey(toSessionId(sessionInfo)); |
|
|
|
} |
|
|
|
} |
|
|
|
|