From 114d2c3c1cee3128329ca1616a76756fde056cba Mon Sep 17 00:00:00 2001 From: Dmytro Skarzhynets Date: Wed, 6 Dec 2023 19:08:28 +0200 Subject: [PATCH] [WIP] Initial implementation --- .../transport/activity/ActivityState.java | 41 ++ .../activity/ActivityStateManager.java | 43 ++ .../activity/ActivityStateManagerImpl.java | 240 +++++++++++ .../activity/ActivityStateReportCallback.java | 41 ++ .../activity/AsyncActivityStateReporter.java | 41 ++ .../service/DefaultTransportService.java | 395 +++++------------- .../service/TransportActivityState.java | 44 ++ 7 files changed, 562 insertions(+), 283 deletions(-) create mode 100644 common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/activity/ActivityState.java create mode 100644 common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/activity/ActivityStateManager.java create mode 100644 common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/activity/ActivityStateManagerImpl.java create mode 100644 common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/activity/ActivityStateReportCallback.java create mode 100644 common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/activity/AsyncActivityStateReporter.java create mode 100644 common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/service/TransportActivityState.java diff --git a/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/activity/ActivityState.java b/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/activity/ActivityState.java new file mode 100644 index 0000000000..09285f9192 --- /dev/null +++ b/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; + +} diff --git a/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/activity/ActivityStateManager.java b/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/activity/ActivityStateManager.java new file mode 100644 index 0000000000..db24c2388d --- /dev/null +++ b/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 { + + void init(); + + void recordActivity(Key key, Supplier newActivityStateSupplier); + + void destroy(); + +} diff --git a/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/activity/ActivityStateManagerImpl.java b/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/activity/ActivityStateManagerImpl.java new file mode 100644 index 0000000000..7d87b02f19 --- /dev/null +++ b/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 implements ActivityStateManager { + + private final ConcurrentMap 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 reporter; + private final long reportingPeriodDurationMillis; + private final String name; + private boolean initialized; + + public ActivityStateManagerImpl(AsyncActivityStateReporter 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 newActivityStateSupplier) { + long newLastRecordedTime = System.currentTimeMillis(); + int capturedPeriodId = currentPeriodId.get(); + lock.readLock().lock(); + SettableFuture> 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 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 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(); + } + } + } + } + +} diff --git a/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/activity/ActivityStateReportCallback.java b/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/activity/ActivityStateReportCallback.java new file mode 100644 index 0000000000..d9d897d4a1 --- /dev/null +++ b/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 { + + void onSuccess(Key key, long reportedTime); + + void onFailure(Key key, Throwable t); + + void onRemove(Key key); + +} diff --git a/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/activity/AsyncActivityStateReporter.java b/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/activity/AsyncActivityStateReporter.java new file mode 100644 index 0000000000..ab9094c252 --- /dev/null +++ b/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 { + + void reportAsync(Key key, State activityState, ActivityStateReportCallback reportCallback); + + void reportAsync(Map activityStates, ActivityStateReportCallback reportCallback); + +} diff --git a/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/service/DefaultTransportService.java b/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/service/DefaultTransportService.java index 0ece919d7b..7769762b84 100644 --- a/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/service/DefaultTransportService.java +++ b/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 sessions = new ConcurrentHashMap<>(); - private final ActivityStateManager activityStateManager; + private ActivityStateManager activityStateManager; private final Map 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 reporter = new AsyncActivityStateReporter<>() { + @Override + public void reportAsync(UUID sessionId, TransportActivityState activityState, ActivityStateReportCallback 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 activityStates, ActivityStateReportCallback reportCallback) { + long expirationTime = System.currentTimeMillis() - sessionInactivityTimeout; + for (Map.Entry 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 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 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() { - @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 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 activityStatesToRemove = new HashSet<>(); - long expirationTime = System.currentTimeMillis() - sessionInactivityTimeout; - for (Map.Entry 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 msgAcknowledgedCallback) { - SessionMetaData sessionMetaData = sessions.get(toSessionId(sessionInfo)); - reportActivityStateToCore(sessionInfo, sessionMetaData, lastActivityTime, msgAcknowledgedCallback); - } - - private void reportActivityStateToCore( - TransportProtos.SessionInfoProto sessionInfo, SessionMetaData sessionMetaData, long lastActivityTime, TransportServiceCallback 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) { diff --git a/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/service/TransportActivityState.java b/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/service/TransportActivityState.java new file mode 100644 index 0000000000..470699ec5d --- /dev/null +++ b/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; + +}