Browse Source

[WIP] Initial implementation

pull/9980/head
Dmytro Skarzhynets 3 years ago
parent
commit
114d2c3c1c
  1. 41
      common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/activity/ActivityState.java
  2. 43
      common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/activity/ActivityStateManager.java
  3. 240
      common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/activity/ActivityStateManagerImpl.java
  4. 41
      common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/activity/ActivityStateReportCallback.java
  5. 41
      common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/activity/AsyncActivityStateReporter.java
  6. 395
      common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/service/DefaultTransportService.java
  7. 44
      common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/service/TransportActivityState.java

41
common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/activity/ActivityState.java

@ -0,0 +1,41 @@
/**
* ThingsBoard, Inc. ("COMPANY") CONFIDENTIAL
*
* Copyright © 2016-2023 ThingsBoard, Inc. All Rights Reserved.
*
* NOTICE: All information contained herein is, and remains
* the property of ThingsBoard, Inc. and its suppliers,
* if any. The intellectual and technical concepts contained
* herein are proprietary to ThingsBoard, Inc.
* and its suppliers and may be covered by U.S. and Foreign Patents,
* patents in process, and are protected by trade secret or copyright law.
*
* Dissemination of this information or reproduction of this material is strictly forbidden
* unless prior written permission is obtained from COMPANY.
*
* Access to the source code contained herein is hereby forbidden to anyone except current COMPANY employees,
* managers or contractors who have executed Confidentiality and Non-disclosure agreements
* explicitly covering such access.
*
* The copyright notice above does not evidence any actual or intended publication
* or disclosure of this source code, which includes
* information that is confidential and/or proprietary, and is a trade secret, of COMPANY.
* ANY REPRODUCTION, MODIFICATION, DISTRIBUTION, PUBLIC PERFORMANCE,
* OR PUBLIC DISPLAY OF OR THROUGH USE OF THIS SOURCE CODE WITHOUT
* THE EXPRESS WRITTEN CONSENT OF COMPANY IS STRICTLY PROHIBITED,
* AND IN VIOLATION OF APPLICABLE LAWS AND INTERNATIONAL TREATIES.
* THE RECEIPT OR POSSESSION OF THIS SOURCE CODE AND/OR RELATED INFORMATION
* DOES NOT CONVEY OR IMPLY ANY RIGHTS TO REPRODUCE, DISCLOSE OR DISTRIBUTE ITS CONTENTS,
* OR TO MANUFACTURE, USE, OR SELL ANYTHING THAT IT MAY DESCRIBE, IN WHOLE OR IN PART.
*/
package org.thingsboard.server.common.transport.activity;
import lombok.Data;
@Data
public class ActivityState {
private volatile long lastRecordedTime;
private volatile long lastReportedTime;
}

43
common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/activity/ActivityStateManager.java

@ -0,0 +1,43 @@
/**
* ThingsBoard, Inc. ("COMPANY") CONFIDENTIAL
*
* Copyright © 2016-2023 ThingsBoard, Inc. All Rights Reserved.
*
* NOTICE: All information contained herein is, and remains
* the property of ThingsBoard, Inc. and its suppliers,
* if any. The intellectual and technical concepts contained
* herein are proprietary to ThingsBoard, Inc.
* and its suppliers and may be covered by U.S. and Foreign Patents,
* patents in process, and are protected by trade secret or copyright law.
*
* Dissemination of this information or reproduction of this material is strictly forbidden
* unless prior written permission is obtained from COMPANY.
*
* Access to the source code contained herein is hereby forbidden to anyone except current COMPANY employees,
* managers or contractors who have executed Confidentiality and Non-disclosure agreements
* explicitly covering such access.
*
* The copyright notice above does not evidence any actual or intended publication
* or disclosure of this source code, which includes
* information that is confidential and/or proprietary, and is a trade secret, of COMPANY.
* ANY REPRODUCTION, MODIFICATION, DISTRIBUTION, PUBLIC PERFORMANCE,
* OR PUBLIC DISPLAY OF OR THROUGH USE OF THIS SOURCE CODE WITHOUT
* THE EXPRESS WRITTEN CONSENT OF COMPANY IS STRICTLY PROHIBITED,
* AND IN VIOLATION OF APPLICABLE LAWS AND INTERNATIONAL TREATIES.
* THE RECEIPT OR POSSESSION OF THIS SOURCE CODE AND/OR RELATED INFORMATION
* DOES NOT CONVEY OR IMPLY ANY RIGHTS TO REPRODUCE, DISCLOSE OR DISTRIBUTE ITS CONTENTS,
* OR TO MANUFACTURE, USE, OR SELL ANYTHING THAT IT MAY DESCRIBE, IN WHOLE OR IN PART.
*/
package org.thingsboard.server.common.transport.activity;
import java.util.function.Supplier;
public interface ActivityStateManager<Key, State extends ActivityState> {
void init();
void recordActivity(Key key, Supplier<State> newActivityStateSupplier);
void destroy();
}

240
common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/activity/ActivityStateManagerImpl.java

@ -0,0 +1,240 @@
/**
* ThingsBoard, Inc. ("COMPANY") CONFIDENTIAL
*
* Copyright © 2016-2023 ThingsBoard, Inc. All Rights Reserved.
*
* NOTICE: All information contained herein is, and remains
* the property of ThingsBoard, Inc. and its suppliers,
* if any. The intellectual and technical concepts contained
* herein are proprietary to ThingsBoard, Inc.
* and its suppliers and may be covered by U.S. and Foreign Patents,
* patents in process, and are protected by trade secret or copyright law.
*
* Dissemination of this information or reproduction of this material is strictly forbidden
* unless prior written permission is obtained from COMPANY.
*
* Access to the source code contained herein is hereby forbidden to anyone except current COMPANY employees,
* managers or contractors who have executed Confidentiality and Non-disclosure agreements
* explicitly covering such access.
*
* The copyright notice above does not evidence any actual or intended publication
* or disclosure of this source code, which includes
* information that is confidential and/or proprietary, and is a trade secret, of COMPANY.
* ANY REPRODUCTION, MODIFICATION, DISTRIBUTION, PUBLIC PERFORMANCE,
* OR PUBLIC DISPLAY OF OR THROUGH USE OF THIS SOURCE CODE WITHOUT
* THE EXPRESS WRITTEN CONSENT OF COMPANY IS STRICTLY PROHIBITED,
* AND IN VIOLATION OF APPLICABLE LAWS AND INTERNATIONAL TREATIES.
* THE RECEIPT OR POSSESSION OF THIS SOURCE CODE AND/OR RELATED INFORMATION
* DOES NOT CONVEY OR IMPLY ANY RIGHTS TO REPRODUCE, DISCLOSE OR DISTRIBUTE ITS CONTENTS,
* OR TO MANUFACTURE, USE, OR SELL ANYTHING THAT IT MAY DESCRIBE, IN WHOLE OR IN PART.
*/
package org.thingsboard.server.common.transport.activity;
import com.google.common.util.concurrent.FutureCallback;
import com.google.common.util.concurrent.Futures;
import com.google.common.util.concurrent.MoreExecutors;
import com.google.common.util.concurrent.SettableFuture;
import lombok.Data;
import lombok.extern.slf4j.Slf4j;
import org.checkerframework.checker.nullness.qual.NonNull;
import org.springframework.data.util.Pair;
import org.thingsboard.common.util.ThingsBoardThreadFactory;
import java.util.Map;
import java.util.Objects;
import java.util.Random;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.ConcurrentMap;
import java.util.concurrent.Executors;
import java.util.concurrent.ScheduledExecutorService;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicInteger;
import java.util.concurrent.locks.ReadWriteLock;
import java.util.concurrent.locks.ReentrantReadWriteLock;
import java.util.function.Supplier;
import java.util.stream.Collectors;
@Slf4j
public final class ActivityStateManagerImpl<Key, State extends ActivityState> implements ActivityStateManager<Key, State> {
private final ConcurrentMap<Key, ActivityStateWrapper> activityStates = new ConcurrentHashMap<>();
@Data
private static class ActivityStateWrapper {
volatile ActivityState activityState;
volatile boolean alreadyBeenReported;
}
private ScheduledExecutorService scheduler;
private final ReadWriteLock lock = new ReentrantReadWriteLock();
private final AtomicInteger currentPeriodId = new AtomicInteger();
private final AsyncActivityStateReporter<Key, State> reporter;
private final long reportingPeriodDurationMillis;
private final String name;
private boolean initialized;
public ActivityStateManagerImpl(AsyncActivityStateReporter<Key, State> reporter, long reportingPeriodDurationMillis, String name) {
this.reporter = Objects.requireNonNull(reporter, "Failed to initialize activity manager: provided reporter is null.");
this.reportingPeriodDurationMillis = reportingPeriodDurationMillis; // TODO: add min/max duration validation
this.name = name == null ? "activity-state-manager" : name;
}
@Override
public synchronized void init() {
if (!initialized) {
initialized = true;
scheduler = Executors.newSingleThreadScheduledExecutor(ThingsBoardThreadFactory.forName(name));
scheduler.scheduleAtFixedRate(this::reportLastEventAndStartNewPeriod, new Random().nextInt((int) reportingPeriodDurationMillis), reportingPeriodDurationMillis, TimeUnit.MILLISECONDS);
}
}
@Override
public void recordActivity(Key activityKey, Supplier<State> newActivityStateSupplier) {
long newLastRecordedTime = System.currentTimeMillis();
int capturedPeriodId = currentPeriodId.get();
lock.readLock().lock();
SettableFuture<Pair<Key, Long>> reportCompletedFuture = SettableFuture.create();
try {
activityStates.compute(activityKey, (key, activityStateWrapper) -> {
if (activityStateWrapper == null) {
State activityState = newActivityStateSupplier.get();
activityState.setLastRecordedTime(newLastRecordedTime);
activityState.setLastReportedTime(0L);
activityStateWrapper = new ActivityStateWrapper();
activityStateWrapper.setActivityState(activityState);
activityStateWrapper.setAlreadyBeenReported(false);
} else {
activityStateWrapper.getActivityState().setLastRecordedTime(newLastRecordedTime);
}
if (activityStateWrapper.isAlreadyBeenReported()) {
return activityStateWrapper;
}
var activityState = activityStateWrapper.getActivityState();
if (activityState.getLastReportedTime() < activityState.getLastRecordedTime()) {
reporter.reportAsync(key, (State) activityState, new ActivityStateReportCallback<>() {
@Override
public void onSuccess(Key key, long reportedTime) {
reportCompletedFuture.set(Pair.of(key, reportedTime));
}
@Override
public void onFailure(Key key, Throwable t) {
reportCompletedFuture.setException(t);
}
@Override
public void onRemove(Key key) {
lock.readLock().lock();
try {
activityStates.remove(key);
} finally {
lock.readLock().unlock();
}
}
});
}
activityStateWrapper.setAlreadyBeenReported(true);
return activityStateWrapper;
});
} finally {
lock.readLock().unlock();
}
Futures.addCallback(reportCompletedFuture, new FutureCallback<>() {
@Override
public void onSuccess(Pair<Key, Long> reportResult) {
lock.readLock().lock();
try {
updateLastReportedTime(reportResult.getFirst(), reportResult.getSecond());
} finally {
lock.readLock().unlock();
}
}
@Override
public void onFailure(@NonNull Throwable t) { // TODO: add failure logging
lock.readLock().lock();
try {
rollbackReportedStatus(activityKey, capturedPeriodId, newLastRecordedTime);
} finally {
lock.readLock().unlock();
}
}
}, MoreExecutors.directExecutor());
}
private void rollbackReportedStatus(Key key, int capturedPeriodId, long newLastRecordedTime) {
activityStates.computeIfPresent(key, (__, activityStateWrapper) -> {
if (capturedPeriodId == currentPeriodId.get() && activityStateWrapper.getActivityState().getLastReportedTime() < newLastRecordedTime) {
activityStateWrapper.setAlreadyBeenReported(false);
}
return activityStateWrapper;
});
}
private void reportLastEventAndStartNewPeriod() {
lock.writeLock().lock();
try {
Map<Key, State> activityStateMap = activityStates.entrySet().stream()
.peek(entry -> entry.getValue().setAlreadyBeenReported(false))
.collect(Collectors.toMap(Map.Entry::getKey, entry -> (State) entry.getValue().getActivityState()));
reporter.reportAsync(activityStateMap, new ActivityStateReportCallback<>() {
@Override
public void onSuccess(Key key, long reportedTime) {
lock.readLock().lock();
try {
updateLastReportedTime(key, reportedTime);
} finally {
lock.readLock().unlock();
}
}
@Override
public void onFailure(Key key, Throwable t) {
log.debug("Failed to report activity state for key [{}]!", key, t);
}
@Override
public void onRemove(Key key) {
lock.readLock().lock();
try {
activityStates.remove(key);
} finally {
lock.readLock().unlock();
}
}
});
} finally {
currentPeriodId.incrementAndGet();
lock.writeLock().unlock();
}
}
private void updateLastReportedTime(Key key, long newLastReportedTime) {
activityStates.computeIfPresent(key, (__, activityStateWrapper) -> {
var activityState = activityStateWrapper.getActivityState();
activityState.setLastReportedTime(Math.max(activityState.getLastReportedTime(), newLastReportedTime));
return activityStateWrapper;
});
}
@Override
public synchronized void destroy() {
if (initialized) {
initialized = false;
if (scheduler != null) {
scheduler.shutdown();
try {
if (scheduler.awaitTermination(10L, TimeUnit.SECONDS)) {
scheduler.shutdownNow();
}
} catch (InterruptedException e) {
scheduler.shutdownNow();
Thread.currentThread().interrupt();
}
}
}
}
}

41
common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/activity/ActivityStateReportCallback.java

@ -0,0 +1,41 @@
/**
* ThingsBoard, Inc. ("COMPANY") CONFIDENTIAL
*
* Copyright © 2016-2023 ThingsBoard, Inc. All Rights Reserved.
*
* NOTICE: All information contained herein is, and remains
* the property of ThingsBoard, Inc. and its suppliers,
* if any. The intellectual and technical concepts contained
* herein are proprietary to ThingsBoard, Inc.
* and its suppliers and may be covered by U.S. and Foreign Patents,
* patents in process, and are protected by trade secret or copyright law.
*
* Dissemination of this information or reproduction of this material is strictly forbidden
* unless prior written permission is obtained from COMPANY.
*
* Access to the source code contained herein is hereby forbidden to anyone except current COMPANY employees,
* managers or contractors who have executed Confidentiality and Non-disclosure agreements
* explicitly covering such access.
*
* The copyright notice above does not evidence any actual or intended publication
* or disclosure of this source code, which includes
* information that is confidential and/or proprietary, and is a trade secret, of COMPANY.
* ANY REPRODUCTION, MODIFICATION, DISTRIBUTION, PUBLIC PERFORMANCE,
* OR PUBLIC DISPLAY OF OR THROUGH USE OF THIS SOURCE CODE WITHOUT
* THE EXPRESS WRITTEN CONSENT OF COMPANY IS STRICTLY PROHIBITED,
* AND IN VIOLATION OF APPLICABLE LAWS AND INTERNATIONAL TREATIES.
* THE RECEIPT OR POSSESSION OF THIS SOURCE CODE AND/OR RELATED INFORMATION
* DOES NOT CONVEY OR IMPLY ANY RIGHTS TO REPRODUCE, DISCLOSE OR DISTRIBUTE ITS CONTENTS,
* OR TO MANUFACTURE, USE, OR SELL ANYTHING THAT IT MAY DESCRIBE, IN WHOLE OR IN PART.
*/
package org.thingsboard.server.common.transport.activity;
public interface ActivityStateReportCallback<Key> {
void onSuccess(Key key, long reportedTime);
void onFailure(Key key, Throwable t);
void onRemove(Key key);
}

41
common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/activity/AsyncActivityStateReporter.java

@ -0,0 +1,41 @@
/**
* ThingsBoard, Inc. ("COMPANY") CONFIDENTIAL
*
* Copyright © 2016-2023 ThingsBoard, Inc. All Rights Reserved.
*
* NOTICE: All information contained herein is, and remains
* the property of ThingsBoard, Inc. and its suppliers,
* if any. The intellectual and technical concepts contained
* herein are proprietary to ThingsBoard, Inc.
* and its suppliers and may be covered by U.S. and Foreign Patents,
* patents in process, and are protected by trade secret or copyright law.
*
* Dissemination of this information or reproduction of this material is strictly forbidden
* unless prior written permission is obtained from COMPANY.
*
* Access to the source code contained herein is hereby forbidden to anyone except current COMPANY employees,
* managers or contractors who have executed Confidentiality and Non-disclosure agreements
* explicitly covering such access.
*
* The copyright notice above does not evidence any actual or intended publication
* or disclosure of this source code, which includes
* information that is confidential and/or proprietary, and is a trade secret, of COMPANY.
* ANY REPRODUCTION, MODIFICATION, DISTRIBUTION, PUBLIC PERFORMANCE,
* OR PUBLIC DISPLAY OF OR THROUGH USE OF THIS SOURCE CODE WITHOUT
* THE EXPRESS WRITTEN CONSENT OF COMPANY IS STRICTLY PROHIBITED,
* AND IN VIOLATION OF APPLICABLE LAWS AND INTERNATIONAL TREATIES.
* THE RECEIPT OR POSSESSION OF THIS SOURCE CODE AND/OR RELATED INFORMATION
* DOES NOT CONVEY OR IMPLY ANY RIGHTS TO REPRODUCE, DISCLOSE OR DISTRIBUTE ITS CONTENTS,
* OR TO MANUFACTURE, USE, OR SELL ANYTHING THAT IT MAY DESCRIBE, IN WHOLE OR IN PART.
*/
package org.thingsboard.server.common.transport.activity;
import java.util.Map;
public interface AsyncActivityStateReporter<Key, State extends ActivityState> {
void reportAsync(Key key, State activityState, ActivityStateReportCallback<Key> reportCallback);
void reportAsync(Map<Key, State> activityStates, ActivityStateReportCallback<Key> reportCallback);
}

395
common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/service/DefaultTransportService.java

@ -72,6 +72,10 @@ import org.thingsboard.server.common.transport.TransportResourceCache;
import org.thingsboard.server.common.transport.TransportService;
import org.thingsboard.server.common.transport.TransportServiceCallback;
import org.thingsboard.server.common.transport.TransportTenantProfileCache;
import org.thingsboard.server.common.transport.activity.ActivityStateManager;
import org.thingsboard.server.common.transport.activity.ActivityStateManagerImpl;
import org.thingsboard.server.common.transport.activity.ActivityStateReportCallback;
import org.thingsboard.server.common.transport.activity.AsyncActivityStateReporter;
import org.thingsboard.server.common.transport.auth.GetOrCreateDeviceFromGatewayResponse;
import org.thingsboard.server.common.transport.auth.TransportDeviceInfo;
import org.thingsboard.server.common.transport.auth.ValidateDeviceCredentialsResponse;
@ -108,14 +112,11 @@ 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.LinkedHashMap;
import java.util.List;
import java.util.Map;
import java.util.Objects;
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;
@ -125,8 +126,6 @@ import java.util.concurrent.Executors;
import java.util.concurrent.ScheduledFuture;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicInteger;
import java.util.concurrent.locks.ReadWriteLock;
import java.util.concurrent.locks.ReentrantReadWriteLock;
import java.util.stream.Collectors;
/**
@ -204,7 +203,7 @@ public class DefaultTransportService implements TransportService {
private ExecutorService mainConsumerExecutor;
public final ConcurrentMap<UUID, SessionMetaData> sessions = new ConcurrentHashMap<>();
private final ActivityStateManager activityStateManager;
private ActivityStateManager<UUID, TransportActivityState> activityStateManager;
private final Map<String, RpcRequestMetadata> toServerRpcPendingMap = new ConcurrentHashMap<>();
private volatile boolean stopped = false;
@ -236,7 +235,6 @@ public class DefaultTransportService implements TransportService {
this.eventPublisher = eventPublisher;
this.notificationRuleProcessor = notificationRuleProcessor;
this.entityLimitsCache = entityLimitsCache;
activityStateManager = new ActivityStateManager();
}
@PostConstruct
@ -245,7 +243,6 @@ public class DefaultTransportService implements TransportService {
this.tbCoreProducerStats = statsFactory.createMessagesStats(StatsType.CORE.getName() + ".producer");
this.transportApiStats = statsFactory.createMessagesStats(StatsType.TRANSPORT.getName() + ".producer");
this.transportCallbackExecutor = ThingsBoardExecutors.newWorkStealingPool(20, getClass());
activityStateManager.init();
this.scheduler.scheduleAtFixedRate(this::invalidateRateLimits, new Random().nextInt((int) sessionReportTimeout), sessionReportTimeout, TimeUnit.MILLISECONDS);
transportApiRequestTemplate = queueProvider.createTransportApiRequestTemplate();
transportApiRequestTemplate.setMessagesStats(transportApiStats);
@ -256,8 +253,108 @@ public class DefaultTransportService implements TransportService {
transportNotificationsConsumer.subscribe(Collections.singleton(tpi));
transportApiRequestTemplate.init();
mainConsumerExecutor = Executors.newSingleThreadExecutor(ThingsBoardThreadFactory.forName("transport-consumer"));
activityStateManager = new ActivityStateManagerImpl<>(reporter, sessionReportTimeout, "transport-activity-state-manager");
activityStateManager.init();
}
// TODO: why sessions and session activities are managed in a separate maps? maybe it is better to manage in a single map
// (eg. we can check if we had sent last activity event when deregistering session)
// TODO: currently if activity event is received between precise session expiration time and reportLastEventAndStartNewPeriod() call
// we will "resurrect" this session meaning that it will be considered alive but in fact it did not perform any activity for more than timeout allows
// TODO: I optimistically set alreadyBeenReported status to true to avoid reporting several activity events if they arrive in rapid succession
// this can cause lost activity updates if first event that got reported failed to enqueue in Kafka and next reporting period is already started
// (maybe having a queue of "first" events will help with this, queue will be cleared once event was successfully queued or later time was reported by event in next period)
// setting alreadyBeenReported status in callbacks ensures that events are reported but can cause several "first" events to be reported
private final AsyncActivityStateReporter<UUID, TransportActivityState> reporter = new AsyncActivityStateReporter<>() {
@Override
public void reportAsync(UUID sessionId, TransportActivityState activityState, ActivityStateReportCallback<UUID> reportCallback) {
SessionMetaData sessionMetaData = sessions.get(sessionId);
TransportProtos.SubscriptionInfoProto subscriptionInfo = TransportProtos.SubscriptionInfoProto.newBuilder()
.setAttributeSubscription(sessionMetaData != null && sessionMetaData.isSubscribedToAttributes())
.setRpcSubscription(sessionMetaData != null && sessionMetaData.isSubscribedToRPC())
.setLastActivityTime(activityState.getLastRecordedTime())
.build();
process(activityState.getSessionInfoProto(), subscriptionInfo, new TransportServiceCallback<>() {
@Override
public void onSuccess(Void msgAcknowledged) {
reportCallback.onSuccess(sessionId, activityState.getLastRecordedTime());
}
@Override
public void onError(Throwable e) {
reportCallback.onFailure(sessionId, e);
}
});
}
@Override
public void reportAsync(Map<UUID, TransportActivityState> activityStates, ActivityStateReportCallback<UUID> reportCallback) {
long expirationTime = System.currentTimeMillis() - sessionInactivityTimeout;
for (Map.Entry<UUID, TransportActivityState> entry : activityStates.entrySet()) {
var sessionId = entry.getKey();
var activityState = entry.getValue();
long lastActivityTime = activityState.getLastRecordedTime();
SessionMetaData sessionMetaData = sessions.get(sessionId);
if (sessionMetaData != null) {
activityState.setSessionInfoProto(sessionMetaData.getSessionInfo());
} else {
reportCallback.onRemove(sessionId);
}
TransportProtos.SessionInfoProto sessionInfo = activityState.getSessionInfoProto();
if (sessionInfo.getGwSessionIdMSB() != 0 && sessionInfo.getGwSessionIdLSB() != 0) {
var gwSessionId = new UUID(sessionInfo.getGwSessionIdMSB(), sessionInfo.getGwSessionIdLSB());
SessionMetaData gwSessionMetaData = sessions.get(gwSessionId);
if (gwSessionMetaData != null && gwSessionMetaData.isOverwriteActivityTime()) {
TransportActivityState gwActivityState = activityStates.get(gwSessionId);
if (gwActivityState != null) {
lastActivityTime = Math.max(gwActivityState.getLastRecordedTime(), lastActivityTime);
}
}
}
if (sessionMetaData != null && lastActivityTime < expirationTime) {
if (log.isDebugEnabled()) {
log.debug("[{}] Session has expired due to last activity time: {}!", sessionId, lastActivityTime);
}
sessions.remove(sessionId);
reportCallback.onRemove(sessionId);
process(sessionInfo, SESSION_EVENT_MSG_CLOSED, null);
sessionMetaData.getListener().onRemoteSessionCloseCommand(sessionId, SESSION_EXPIRED_NOTIFICATION_PROTO);
} else if (activityState.getLastReportedTime() < lastActivityTime) {
long finalLastActivityTime = lastActivityTime;
reportActivityStateToCore(sessionInfo, sessionMetaData, lastActivityTime, new TransportServiceCallback<>() {
@Override
public void onSuccess(Void msgAcknowledged) {
reportCallback.onSuccess(sessionId, finalLastActivityTime);
}
@Override
public void onError(Throwable e) {
reportCallback.onFailure(sessionId, e);
}
});
}
}
}
private void reportActivityStateToCore(
TransportProtos.SessionInfoProto sessionInfo, SessionMetaData sessionMetaData, long lastActivityTime, TransportServiceCallback<Void> callback
) {
TransportProtos.SubscriptionInfoProto subscriptionInfo = TransportProtos.SubscriptionInfoProto.newBuilder()
.setAttributeSubscription(sessionMetaData != null && sessionMetaData.isSubscribedToAttributes())
.setRpcSubscription(sessionMetaData != null && sessionMetaData.isSubscribedToRPC())
.setLastActivityTime(lastActivityTime)
.build();
process(sessionInfo, subscriptionInfo, callback);
}
};
@AfterStartUp(order = AfterStartUp.TRANSPORT_SERVICE)
public void start() {
mainConsumerExecutor.execute(() -> {
@ -309,6 +406,9 @@ public class DefaultTransportService implements TransportService {
if (transportApiRequestTemplate != null) {
transportApiRequestTemplate.stop();
}
if (activityStateManager != null) {
activityStateManager.destroy();
}
}
@Override
@ -788,283 +888,12 @@ public class DefaultTransportService implements TransportService {
}
private void reportActivityInternal(TransportProtos.SessionInfoProto sessionInfo) {
activityStateManager.recordActivity(sessionInfo);
}
/*
private void checkInactivityAndReportActivity() {
long expTime = System.currentTimeMillis() - sessionInactivityTimeout;
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(gwActivityData.getLastActivityTime(), lastActivityTime);
}
}
if (lastActivityTime < expTime) {
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, SESSION_EVENT_MSG_CLOSED, null);
sessionMD.getListener().onRemoteSessionCloseCommand(uuid, SESSION_EXPIRED_NOTIFICATION_PROTO);
}
} else {
if (lastActivityTime > sessionAD.getLastReportedActivityTime()) {
final long lastActivityTimeFinal = lastActivityTime;
process(sessionInfo, TransportProtos.SubscriptionInfoProto.newBuilder()
.setAttributeSubscription(sessionMD != null && sessionMD.isSubscribedToAttributes())
.setRpcSubscription(sessionMD != null && sessionMD.isSubscribedToRPC())
.setLastActivityTime(lastActivityTime).build(), new TransportServiceCallback<Void>() {
@Override
public void onSuccess(Void msg) {
sessionAD.setLastReportedActivityTime(lastActivityTimeFinal);
}
@Override
public void onError(Throwable e) {
log.warn("[{}] Failed to report last activity time", uuid, e);
}
});
}
}
activityStateManager.recordActivity(toSessionId(sessionInfo), () -> {
var activityState = new TransportActivityState();
activityState.setSessionInfoProto(sessionInfo);
return activityState;
});
// Removes all closed or short-lived sessions.
sessionsToRemove.forEach(sessionsActivity::remove);
}
*/
// TODO: how can I get access to sessions if activity management logic is implemented as a separate service (class)
// TODO: why sessions and session activities are managed in a separate maps? maybe it is better to manage in a single map
// (eg. we can check if we had sent last activity event when deregistering session)
// TODO: currently if activity event is received between precise session expiration time and reportLastEventAndStartNewPeriod() call
// we will "resurrect" this session meaning that
// TODO: I optimistically set alreadyBeenReported status to true to avoid reporting several activity events if they arrive in rapid succession
// this can cause lost activity updates if first event that got reported failed to enqueue in Kafka and next reporting period is already started
// (maybe having a queue of "first" events will help with this, queue will be cleared once event was successfully queued or later time was reported by event in next period)
// setting alreadyBeenReported status in callbacks ensures that events are reported but can cause several "first" events to be reported
// TODO: currently activity states are reported on per session basis, but one physical device can have several sessions
// maybe it is a good idea to manage activity state on per device basis
private final class ActivityStateManager {
private final ConcurrentMap<UUID, ActivityState> sessionActivityStates = new ConcurrentHashMap<>();
@lombok.Value
// TODO: I chose to have immutable objects for concurrency considerations,
// but it can possibly create a performance issue with creating large amounts of new objects
// maybe mutable objects with volatile fields is better
private class ActivityState {
TransportProtos.SessionInfoProto sessionInfo;
long lastActivityTime;
boolean alreadyBeenReported;
long lastReportedTime;
ActivityState withSessionInfo(TransportProtos.SessionInfoProto sessionInfo) {
return new ActivityState(sessionInfo, getLastActivityTime(), isAlreadyBeenReported(), getLastReportedTime());
}
ActivityState withLastActivityTime(long lastActivityTime) {
return new ActivityState(getSessionInfo(), lastActivityTime, isAlreadyBeenReported(), getLastReportedTime());
}
ActivityState withReportedStatus(boolean hasAlreadyBeenReported) {
return new ActivityState(getSessionInfo(), getLastActivityTime(), hasAlreadyBeenReported, getLastReportedTime());
}
ActivityState withLastReportedTime(long lastReportedTime) {
return new ActivityState(getSessionInfo(), getLastActivityTime(), isAlreadyBeenReported(), lastReportedTime);
}
}
// TODO: maybe CAS-based synchronization policy will perform better?
// locks can put threads to sleep which introduces scheduling overhead (it is up to JVM to device if a thread will go to sleep or will be spin-waiting)
// activity management tasks are short lived so scheduling overhead may be significant, so CAS can be better since it does not put threads to sleep (always spin-waiting)
// Compare ReadWriteLock to CAS
private final ReadWriteLock lock = new ReentrantReadWriteLock();
private final AtomicInteger currentPeriodId = new AtomicInteger();
private void init() {
scheduler.scheduleAtFixedRate(this::reportLastEventAndStartNewPeriod, new Random().nextInt((int) sessionReportTimeout), sessionReportTimeout, TimeUnit.MILLISECONDS);
}
private void recordActivity(TransportProtos.SessionInfoProto sessionInfo) {
long newLastActivityTime = System.currentTimeMillis();
var sessionId = toSessionId(sessionInfo);
lock.readLock().lock();
try {
sessionActivityStates.compute(sessionId, (id, currentActivityState) -> {
log.info("------------------------------------------------------------------------------------------------------------------------");
log.info("Record activity: entered compute! Session id: [{}]", sessionId);
var activityState = Objects.requireNonNullElseGet(currentActivityState, () -> new ActivityState(sessionInfo, newLastActivityTime, false, 0L));
if (activityState.isAlreadyBeenReported()) { // update the last activity time
log.info("------------------------------------------------------------------------------------------------------------------------");
return activityState.withLastActivityTime(newLastActivityTime);
}
int capturedPeriodId = currentPeriodId.get();
reportActivityStateToCore(activityState.getSessionInfo(), newLastActivityTime, new TransportServiceCallback<>() {
@Override
public void onSuccess(Void msgAcknowledged) {
// Success, the optimistic assumption was correct: update last reported time if it is newer
log.info("Record activity: successful callback received! Session id: [{}]", sessionId);
lock.readLock().lock();
try {
updateLastReportedTime(sessionId, newLastActivityTime);
} finally {
lock.readLock().unlock();
}
}
@Override
public void onError(Throwable e) {
// Error: revert alreadyBeenReported to false
log.info("Record activity: failed callback received! Session id: [{}]", sessionId);
log.debug("[{}] Failed to report last activity time!", sessionId, e);
lock.readLock().lock();
try {
sessionActivityStates.computeIfPresent(sessionId, (__, activityState) -> {
boolean updatedReportedStatus = activityState.isAlreadyBeenReported();
if (capturedPeriodId == currentPeriodId.get() && activityState.getLastReportedTime() < newLastActivityTime) {
updatedReportedStatus = false;
}
return activityState.withReportedStatus(updatedReportedStatus);
});
} finally {
lock.readLock().unlock();
}
}
});
log.info("------------------------------------------------------------------------------------------------------------------------");
// Optimistically set hasAlreadyBeenReported to true
return new ActivityState(sessionInfo, newLastActivityTime, true, activityState.getLastReportedTime());
});
} finally {
lock.readLock().unlock();
}
}
private void reportLastEventAndStartNewPeriod() {
lock.writeLock().lock();
// log.info("------------------------------------------------------------------------------------------------------------------------");
// log.info("Period change start! Current period id: [{}], activity states: {}", currentPeriodId.get(), sessionActivityStates.keySet());
try {
Set<UUID> activityStatesToRemove = new HashSet<>();
long expirationTime = System.currentTimeMillis() - sessionInactivityTimeout;
for (Map.Entry<UUID, ActivityState> entry : sessionActivityStates.entrySet()) {
var sessionId = entry.getKey();
var activityState = entry.getValue();
long lastActivityTime = activityState.getLastActivityTime();
// TODO: why this update is needed? if gateway session is present but this session is already gone we will not check if gateway's last activity time is greater than this one
SessionMetaData sessionMetaData = sessions.get(sessionId);
if (sessionMetaData != null) {
activityState = activityState.withSessionInfo(sessionMetaData.getSessionInfo());
entry.setValue(activityState);
} else {
// log.info("Removing activity state! Session id: [{}]", sessionId);
activityStatesToRemove.add(sessionId);
}
// TODO: ask about this part end
TransportProtos.SessionInfoProto sessionInfo = activityState.getSessionInfo();
if (sessionInfo.getGwSessionIdMSB() != 0 && sessionInfo.getGwSessionIdLSB() != 0) {
var gwSessionId = new UUID(sessionInfo.getGwSessionIdMSB(), sessionInfo.getGwSessionIdLSB());
SessionMetaData gwSessionMetaData = sessions.get(gwSessionId);
if (gwSessionMetaData != null && gwSessionMetaData.isOverwriteActivityTime()) {
ActivityState gwActivityState = sessionActivityStates.get(gwSessionId);
lastActivityTime = Math.max(gwActivityState.getLastActivityTime(), lastActivityTime);
}
}
if (sessionMetaData != null && lastActivityTime < expirationTime) {
// log.info("Session has expired! Session id: [{}]", sessionId);
if (log.isDebugEnabled()) {
log.debug("[{}] Session has expired due to last activity time: {}!", sessionId, lastActivityTime);
}
sessions.remove(sessionId);
activityStatesToRemove.add(sessionId);
process(sessionInfo, SESSION_EVENT_MSG_CLOSED, null);
sessionMetaData.getListener().onRemoteSessionCloseCommand(sessionId, SESSION_EXPIRED_NOTIFICATION_PROTO);
} else if (activityState.getLastReportedTime() < lastActivityTime) {
long finalLastActivityTime = lastActivityTime;
reportActivityStateToCore(sessionInfo, sessionMetaData, lastActivityTime, new TransportServiceCallback<>() {
@Override
public void onSuccess(Void msgAcknowledged) {
lock.readLock().lock();
try {
updateLastReportedTime(sessionId, finalLastActivityTime);
} finally {
lock.readLock().unlock();
}
}
@Override
public void onError(Throwable e) {
log.debug("[{}] Failed to report last activity time!", sessionId, e);
}
});
}
entry.setValue(activityState.withReportedStatus(false));
}
activityStatesToRemove.forEach(sessionActivityStates::remove);
} finally {
currentPeriodId.incrementAndGet();
// log.info("Period change end! Current period id: [{}], activity states {}", currentPeriodId.get(), sessionActivityStates.keySet());
// log.info("------------------------------------------------------------------------------------------------------------------------");
lock.writeLock().unlock();
}
}
private void reportActivityStateToCore(TransportProtos.SessionInfoProto sessionInfo, long lastActivityTime, TransportServiceCallback<Void> msgAcknowledgedCallback) {
SessionMetaData sessionMetaData = sessions.get(toSessionId(sessionInfo));
reportActivityStateToCore(sessionInfo, sessionMetaData, lastActivityTime, msgAcknowledgedCallback);
}
private void reportActivityStateToCore(
TransportProtos.SessionInfoProto sessionInfo, SessionMetaData sessionMetaData, long lastActivityTime, TransportServiceCallback<Void> msgAcknowledgedCallback
) {
log.info("Reporting activity state to core! Session id: [{}], last activity time: [{}], current period id: [{}]",
toSessionId(sessionInfo), lastActivityTime, currentPeriodId.get());
TransportProtos.SubscriptionInfoProto subscriptionInfo = TransportProtos.SubscriptionInfoProto.newBuilder()
.setAttributeSubscription(sessionMetaData != null && sessionMetaData.isSubscribedToAttributes())
.setRpcSubscription(sessionMetaData != null && sessionMetaData.isSubscribedToRPC())
.setLastActivityTime(lastActivityTime)
.build();
process(sessionInfo, subscriptionInfo, msgAcknowledgedCallback);
}
private void updateLastReportedTime(UUID sessionId, long lastActivityTime) {
sessionActivityStates.computeIfPresent(sessionId, (__, currentActivityState) -> {
long updatedLastReportedTime = Math.max(currentActivityState.getLastReportedTime(), lastActivityTime);
return currentActivityState.withLastReportedTime(updatedLastReportedTime);
});
}
}
@Override
public void lifecycleEvent(TenantId tenantId, DeviceId deviceId, ComponentLifecycleEvent eventType, boolean success, Throwable error) {

44
common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/service/TransportActivityState.java

@ -0,0 +1,44 @@
/**
* ThingsBoard, Inc. ("COMPANY") CONFIDENTIAL
*
* Copyright © 2016-2023 ThingsBoard, Inc. All Rights Reserved.
*
* NOTICE: All information contained herein is, and remains
* the property of ThingsBoard, Inc. and its suppliers,
* if any. The intellectual and technical concepts contained
* herein are proprietary to ThingsBoard, Inc.
* and its suppliers and may be covered by U.S. and Foreign Patents,
* patents in process, and are protected by trade secret or copyright law.
*
* Dissemination of this information or reproduction of this material is strictly forbidden
* unless prior written permission is obtained from COMPANY.
*
* Access to the source code contained herein is hereby forbidden to anyone except current COMPANY employees,
* managers or contractors who have executed Confidentiality and Non-disclosure agreements
* explicitly covering such access.
*
* The copyright notice above does not evidence any actual or intended publication
* or disclosure of this source code, which includes
* information that is confidential and/or proprietary, and is a trade secret, of COMPANY.
* ANY REPRODUCTION, MODIFICATION, DISTRIBUTION, PUBLIC PERFORMANCE,
* OR PUBLIC DISPLAY OF OR THROUGH USE OF THIS SOURCE CODE WITHOUT
* THE EXPRESS WRITTEN CONSENT OF COMPANY IS STRICTLY PROHIBITED,
* AND IN VIOLATION OF APPLICABLE LAWS AND INTERNATIONAL TREATIES.
* THE RECEIPT OR POSSESSION OF THIS SOURCE CODE AND/OR RELATED INFORMATION
* DOES NOT CONVEY OR IMPLY ANY RIGHTS TO REPRODUCE, DISCLOSE OR DISTRIBUTE ITS CONTENTS,
* OR TO MANUFACTURE, USE, OR SELL ANYTHING THAT IT MAY DESCRIBE, IN WHOLE OR IN PART.
*/
package org.thingsboard.server.common.transport.service;
import lombok.Data;
import lombok.EqualsAndHashCode;
import org.thingsboard.server.common.transport.activity.ActivityState;
import org.thingsboard.server.gen.transport.TransportProtos;
@Data
@EqualsAndHashCode(callSuper = true)
public class TransportActivityState extends ActivityState {
private volatile TransportProtos.SessionInfoProto sessionInfoProto;
}
Loading…
Cancel
Save