From f93a49be6d4a795726e57132d4a197eb44f5aa7e Mon Sep 17 00:00:00 2001 From: Dmytro Skarzhynets Date: Thu, 14 Dec 2023 23:14:44 +0200 Subject: [PATCH] [WIP] Refactoring after review --- .../AllEventsIntegrationActivityManager.java | 124 -------------- ...dLastEventsIntegrationActivityManager.java | 157 ----------------- ...stEventOnlyIntegrationActivityManager.java | 160 ------------------ ...stEventOnlyIntegrationActivityManager.java | 105 ------------ .../src/main/resources/thingsboard.yml | 13 +- .../activity/AbstractActivityManager.java | 103 +++++++++-- .../transport/activity/ActivityManager.java | 16 +- .../transport/activity/ActivityState.java | 5 +- .../activity/strategy/ActivityStrategy.java | 11 ++ .../strategy/ActivityStrategyType.java | 20 +++ .../strategy/AllEventsActivityStrategy.java | 17 ++ .../FirstAndLastEventActivityStrategy.java | 25 +++ .../strategy/FirstEventActivityStrategy.java | 25 +++ .../strategy/LastEventActivityStrategy.java | 17 ++ 14 files changed, 223 insertions(+), 575 deletions(-) delete mode 100644 application/src/main/java/org/thingsboard/server/service/integration/activity/AllEventsIntegrationActivityManager.java delete mode 100644 application/src/main/java/org/thingsboard/server/service/integration/activity/FirstAndLastEventsIntegrationActivityManager.java delete mode 100644 application/src/main/java/org/thingsboard/server/service/integration/activity/FirstEventOnlyIntegrationActivityManager.java delete mode 100644 application/src/main/java/org/thingsboard/server/service/integration/activity/LastEventOnlyIntegrationActivityManager.java create mode 100644 common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/activity/strategy/ActivityStrategy.java create mode 100644 common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/activity/strategy/ActivityStrategyType.java create mode 100644 common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/activity/strategy/AllEventsActivityStrategy.java create mode 100644 common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/activity/strategy/FirstAndLastEventActivityStrategy.java create mode 100644 common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/activity/strategy/FirstEventActivityStrategy.java create mode 100644 common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/activity/strategy/LastEventActivityStrategy.java diff --git a/application/src/main/java/org/thingsboard/server/service/integration/activity/AllEventsIntegrationActivityManager.java b/application/src/main/java/org/thingsboard/server/service/integration/activity/AllEventsIntegrationActivityManager.java deleted file mode 100644 index f2edb20a7f..0000000000 --- a/application/src/main/java/org/thingsboard/server/service/integration/activity/AllEventsIntegrationActivityManager.java +++ /dev/null @@ -1,124 +0,0 @@ -/** - * 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.service.integration.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.extern.slf4j.Slf4j; -import org.checkerframework.checker.nullness.qual.NonNull; -import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty; -import org.springframework.data.util.Pair; -import org.springframework.stereotype.Component; -import org.thingsboard.server.common.transport.activity.AbstractActivityManager; -import org.thingsboard.server.common.transport.activity.ActivityReportCallback; -import org.thingsboard.server.common.transport.activity.ActivityState; -import org.thingsboard.server.queue.util.TbCoreComponent; - -import java.util.Map; -import java.util.concurrent.ConcurrentHashMap; -import java.util.concurrent.ConcurrentMap; -import java.util.function.Supplier; - -@Slf4j -@Component -@TbCoreComponent -@ConditionalOnProperty(prefix = "integrations.activity", value = "reporting_strategy", havingValue = "all") -public class AllEventsIntegrationActivityManager extends AbstractActivityManager { - - private final ConcurrentMap states = new ConcurrentHashMap<>(); - - @Override - protected void doOnActivity(IntegrationActivityKey activityKey, Supplier newStateSupplier) { - long newLastRecordedTime = System.currentTimeMillis(); - SettableFuture> reportCompletedFuture = SettableFuture.create(); - states.compute(activityKey, (key, activityState) -> { - if (activityState == null) { - activityState = newStateSupplier.get(); - } - if (activityState.getLastRecordedTime() < newLastRecordedTime) { - activityState.setLastRecordedTime(newLastRecordedTime); - } - if (activityState.getLastReportedTime() < activityState.getLastRecordedTime()) { - log.debug("[{}][{}] Going to report activity event for device with id: [{}].", - activityKey.getTenantId().getId(), name, activityKey.getDeviceId().getId()); - reporter.report(key, activityState.getLastRecordedTime(), activityState, new ActivityReportCallback<>() { - @Override - public void onSuccess(IntegrationActivityKey key, long reportedTime) { - reportCompletedFuture.set(Pair.of(key, reportedTime)); - } - - @Override - public void onFailure(IntegrationActivityKey key, Throwable t) { - reportCompletedFuture.setException(t); - } - }); - } - return activityState; - }); - Futures.addCallback(reportCompletedFuture, new FutureCallback<>() { - @Override - public void onSuccess(Pair reportResult) { - updateLastReportedTime(reportResult.getFirst(), reportResult.getSecond()); - } - - @Override - public void onFailure(@NonNull Throwable t) { - log.debug("[{}][{}] Failed to report activity event for device with id: [{}].", - name, activityKey.getTenantId().getId(), activityKey.getDeviceId().getId()); - } - }, MoreExecutors.directExecutor()); - } - - private void updateLastReportedTime(IntegrationActivityKey key, long newLastReportedTime) { - states.computeIfPresent(key, (__, activityState) -> { - activityState.setLastReportedTime(Math.max(activityState.getLastReportedTime(), newLastReportedTime)); - return activityState; - }); - } - - @Override - protected void doOnReportingPeriodEnd() { - for (Map.Entry entry : states.entrySet()) { - var activityKey = entry.getKey(); - var activityState = entry.getValue(); - // if there were no activities during the reporting period, we should remove the entry to prevent memory leaks - long expirationTime = System.currentTimeMillis() - reportingPeriodMillis; - if (activityState.getLastRecordedTime() < expirationTime) { - log.debug("[{}][{}] No activity events were received during reporting period for device with id: [{}]. Going to remove activity state.", - activityKey.getTenantId().getId(), name, activityKey.getDeviceId().getId()); - states.remove(activityKey); - } - } - } - -} diff --git a/application/src/main/java/org/thingsboard/server/service/integration/activity/FirstAndLastEventsIntegrationActivityManager.java b/application/src/main/java/org/thingsboard/server/service/integration/activity/FirstAndLastEventsIntegrationActivityManager.java deleted file mode 100644 index 1c2dba809e..0000000000 --- a/application/src/main/java/org/thingsboard/server/service/integration/activity/FirstAndLastEventsIntegrationActivityManager.java +++ /dev/null @@ -1,157 +0,0 @@ -/** - * 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.service.integration.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.boot.autoconfigure.condition.ConditionalOnProperty; -import org.springframework.data.util.Pair; -import org.springframework.stereotype.Component; -import org.thingsboard.server.common.transport.activity.AbstractActivityManager; -import org.thingsboard.server.common.transport.activity.ActivityReportCallback; -import org.thingsboard.server.common.transport.activity.ActivityState; -import org.thingsboard.server.queue.util.TbCoreComponent; - -import java.util.Map; -import java.util.concurrent.ConcurrentHashMap; -import java.util.concurrent.ConcurrentMap; -import java.util.function.Supplier; - -@Slf4j -@Component -@TbCoreComponent -@ConditionalOnProperty(prefix = "integrations.activity", value = "reporting_strategy", havingValue = "first-and-last") -public class FirstAndLastEventsIntegrationActivityManager extends AbstractActivityManager { - - private final ConcurrentMap states = new ConcurrentHashMap<>(); - - @Data - private static class ActivityStateWrapper { - - volatile ActivityState state; - volatile boolean alreadyBeenReported; - - } - - @Override - protected void doOnActivity(IntegrationActivityKey activityKey, Supplier newStateSupplier) { - long newLastRecordedTime = System.currentTimeMillis(); - SettableFuture> reportCompletedFuture = SettableFuture.create(); - states.compute(activityKey, (key, activityStateWrapper) -> { - if (activityStateWrapper == null) { - activityStateWrapper = new ActivityStateWrapper(); - activityStateWrapper.setState(newStateSupplier.get()); - } - var activityState = activityStateWrapper.getState(); - if (activityState.getLastRecordedTime() < newLastRecordedTime) { - activityState.setLastRecordedTime(newLastRecordedTime); - } - if (activityStateWrapper.isAlreadyBeenReported()) { - return activityStateWrapper; - } - if (activityState.getLastReportedTime() < activityState.getLastRecordedTime()) { - log.debug("[{}][{}] Going to report first activity event for device with id: [{}].", - activityKey.getTenantId().getId(), name, activityKey.getDeviceId().getId()); - reporter.report(key, activityState.getLastRecordedTime(), activityState, new ActivityReportCallback<>() { - @Override - public void onSuccess(IntegrationActivityKey key, long reportedTime) { - reportCompletedFuture.set(Pair.of(key, reportedTime)); - } - - @Override - public void onFailure(IntegrationActivityKey key, Throwable t) { - reportCompletedFuture.setException(t); - } - }); - } - activityStateWrapper.setAlreadyBeenReported(true); - return activityStateWrapper; - }); - Futures.addCallback(reportCompletedFuture, new FutureCallback<>() { - @Override - public void onSuccess(Pair reportResult) { - updateLastReportedTime(reportResult.getFirst(), reportResult.getSecond()); - } - - @Override - public void onFailure(@NonNull Throwable t) { - log.debug("[{}][{}] Failed to report first activity event for device with id: [{}].", - activityKey.getTenantId().getId(), name, activityKey.getDeviceId().getId()); - } - }, MoreExecutors.directExecutor()); - } - - @Override - protected void doOnReportingPeriodEnd() { - for (Map.Entry entry : states.entrySet()) { - var activityKey = entry.getKey(); - var activityStateWrapper = entry.getValue(); - var activityState = activityStateWrapper.getState(); - long lastRecordedTime = activityState.getLastRecordedTime(); - // if there were no activities during the reporting period, we should remove the entry to prevent memory leaks - if (!activityStateWrapper.isAlreadyBeenReported()) { - log.debug("[{}][{}] No activity events were received during reporting period for device with id: [{}]. Going to remove activity state.", - activityKey.getTenantId().getId(), name, activityKey.getDeviceId().getId()); - states.remove(activityKey); - } - if (activityState.getLastReportedTime() < lastRecordedTime) { - log.debug("[{}][{}] Going to report last activity event for device with id: [{}].", activityKey.getTenantId().getId(), name, activityKey.getDeviceId().getId()); - reporter.report(activityKey, lastRecordedTime, activityState, new ActivityReportCallback<>() { - @Override - public void onSuccess(IntegrationActivityKey key, long newLastReportedTime) { - updateLastReportedTime(key, newLastReportedTime); - } - - @Override - public void onFailure(IntegrationActivityKey key, Throwable t) { - log.debug("[{}][{}] Failed to report last activity event in a period for device with id: [{}].", - activityKey.getTenantId().getId(), name, activityKey.getDeviceId().getId()); - } - }); - } - activityStateWrapper.setAlreadyBeenReported(false); - } - } - - private void updateLastReportedTime(IntegrationActivityKey key, long newLastReportedTime) { - states.computeIfPresent(key, (__, activityStateWrapper) -> { - var activityState = activityStateWrapper.getState(); - activityState.setLastReportedTime(Math.max(activityState.getLastReportedTime(), newLastReportedTime)); - return activityStateWrapper; - }); - } - -} diff --git a/application/src/main/java/org/thingsboard/server/service/integration/activity/FirstEventOnlyIntegrationActivityManager.java b/application/src/main/java/org/thingsboard/server/service/integration/activity/FirstEventOnlyIntegrationActivityManager.java deleted file mode 100644 index 24f9a8f388..0000000000 --- a/application/src/main/java/org/thingsboard/server/service/integration/activity/FirstEventOnlyIntegrationActivityManager.java +++ /dev/null @@ -1,160 +0,0 @@ -/** - * 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.service.integration.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.boot.autoconfigure.condition.ConditionalOnProperty; -import org.springframework.data.util.Pair; -import org.springframework.stereotype.Component; -import org.thingsboard.server.common.transport.activity.AbstractActivityManager; -import org.thingsboard.server.common.transport.activity.ActivityReportCallback; -import org.thingsboard.server.common.transport.activity.ActivityState; -import org.thingsboard.server.queue.util.TbCoreComponent; - -import java.util.Map; -import java.util.concurrent.ConcurrentHashMap; -import java.util.concurrent.ConcurrentMap; -import java.util.function.Supplier; - -@Slf4j -@Component -@TbCoreComponent -@ConditionalOnProperty(prefix = "integrations.activity", value = "reporting_strategy", havingValue = "first") -public class FirstEventOnlyIntegrationActivityManager extends AbstractActivityManager { - - private final ConcurrentMap states = new ConcurrentHashMap<>(); - - @Data - private static class ActivityStateWrapper { - - volatile ActivityState state; - volatile boolean alreadyBeenReported; - - } - - @Override - protected void doOnActivity(IntegrationActivityKey activityKey, Supplier newStateSupplier) { - long newLastRecordedTime = System.currentTimeMillis(); - SettableFuture> reportCompletedFuture = SettableFuture.create(); - states.compute(activityKey, (key, activityStateWrapper) -> { - if (activityStateWrapper == null) { - activityStateWrapper = new ActivityStateWrapper(); - activityStateWrapper.setState(newStateSupplier.get()); - } - var activityState = activityStateWrapper.getState(); - if (activityState.getLastRecordedTime() < newLastRecordedTime) { - activityState.setLastRecordedTime(newLastRecordedTime); - } - if (activityStateWrapper.isAlreadyBeenReported()) { - return activityStateWrapper; - } - if (activityState.getLastReportedTime() < activityState.getLastRecordedTime()) { - log.debug("[{}][{}] Going to report first activity event for device with id: [{}].", - activityKey.getTenantId().getId(), name, activityKey.getDeviceId().getId()); - reporter.report(key, activityState.getLastRecordedTime(), activityState, new ActivityReportCallback<>() { - @Override - public void onSuccess(IntegrationActivityKey key, long reportedTime) { - reportCompletedFuture.set(Pair.of(key, reportedTime)); - } - - @Override - public void onFailure(IntegrationActivityKey key, Throwable t) { - reportCompletedFuture.setException(t); - } - }); - } - activityStateWrapper.setAlreadyBeenReported(true); - return activityStateWrapper; - }); - Futures.addCallback(reportCompletedFuture, new FutureCallback<>() { - @Override - public void onSuccess(Pair reportResult) { - updateLastReportedTime(reportResult.getFirst(), reportResult.getSecond()); - } - - @Override - public void onFailure(@NonNull Throwable t) { - log.debug("[{}][{}] Failed to report first activity event for device with id: [{}].", - name, activityKey.getTenantId().getId(), activityKey.getDeviceId().getId()); - } - }, MoreExecutors.directExecutor()); - } - - @Override - protected void doOnReportingPeriodEnd() { - for (Map.Entry entry : states.entrySet()) { - var activityStateWrapper = entry.getValue(); - if (activityStateWrapper.isAlreadyBeenReported()) { - activityStateWrapper.setAlreadyBeenReported(false); - continue; - } - var activityKey = entry.getKey(); - var activityState = activityStateWrapper.getState(); - long lastRecordedTime = activityState.getLastRecordedTime(); - // if there were no activities during the reporting period, we should remove the entry to prevent memory leaks - log.debug("[{}][{}] No activity events were received during reporting period for device with id: [{}]. Going to remove activity state.", - activityKey.getTenantId().getId(), name, activityKey.getDeviceId().getId()); - states.remove(activityKey); - // report leftover events - if (activityState.getLastReportedTime() < lastRecordedTime) { - log.debug("[{}][{}] Going to report leftover activity event for device with id: [{}].", activityKey.getTenantId().getId(), name, activityKey.getDeviceId().getId()); - reporter.report(activityKey, lastRecordedTime, activityState, new ActivityReportCallback<>() { - @Override - public void onSuccess(IntegrationActivityKey key, long reportedTime) { - updateLastReportedTime(key, reportedTime); // just in case the same key was added again - } - - @Override - public void onFailure(IntegrationActivityKey key, Throwable t) { - log.debug("[{}][{}] Failed to report last activity event in a period for device with id: [{}].", - activityKey.getTenantId().getId(), name, activityKey.getDeviceId().getId()); - } - }); - } - activityStateWrapper.setAlreadyBeenReported(false); - } - } - - private void updateLastReportedTime(IntegrationActivityKey key, long newLastReportedTime) { - states.computeIfPresent(key, (__, activityStateWrapper) -> { - var activityState = activityStateWrapper.getState(); - activityState.setLastReportedTime(Math.max(activityState.getLastReportedTime(), newLastReportedTime)); - return activityStateWrapper; - }); - } - -} diff --git a/application/src/main/java/org/thingsboard/server/service/integration/activity/LastEventOnlyIntegrationActivityManager.java b/application/src/main/java/org/thingsboard/server/service/integration/activity/LastEventOnlyIntegrationActivityManager.java deleted file mode 100644 index 16699b2e3c..0000000000 --- a/application/src/main/java/org/thingsboard/server/service/integration/activity/LastEventOnlyIntegrationActivityManager.java +++ /dev/null @@ -1,105 +0,0 @@ -/** - * 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.service.integration.activity; - -import lombok.extern.slf4j.Slf4j; -import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty; -import org.springframework.stereotype.Component; -import org.thingsboard.server.common.transport.activity.AbstractActivityManager; -import org.thingsboard.server.common.transport.activity.ActivityReportCallback; -import org.thingsboard.server.common.transport.activity.ActivityState; -import org.thingsboard.server.queue.util.TbCoreComponent; - -import java.util.Map; -import java.util.concurrent.ConcurrentHashMap; -import java.util.concurrent.ConcurrentMap; -import java.util.function.Supplier; - -@Slf4j -@Component -@TbCoreComponent -@ConditionalOnProperty(prefix = "integrations.activity", value = "reporting_strategy", havingValue = "last") -public class LastEventOnlyIntegrationActivityManager extends AbstractActivityManager { - - private final ConcurrentMap states = new ConcurrentHashMap<>(); - - @Override - protected void doOnActivity(IntegrationActivityKey activityKey, Supplier newStateSupplier) { - long newLastRecordedTime = System.currentTimeMillis(); - states.compute(activityKey, (__, activityState) -> { - if (activityState == null) { - activityState = newStateSupplier.get(); - } - if (activityState.getLastRecordedTime() < newLastRecordedTime) { - activityState.setLastRecordedTime(newLastRecordedTime); - } - return activityState; - }); - } - - @Override - public void doOnReportingPeriodEnd() { - for (Map.Entry entry : states.entrySet()) { - var activityKey = entry.getKey(); - var activityState = entry.getValue(); - long lastRecordedTime = activityState.getLastRecordedTime(); - // if there were no activities during the reporting period, we should remove the entry to prevent memory leaks - long expirationTime = System.currentTimeMillis() - reportingPeriodMillis; - if (lastRecordedTime < expirationTime) { - log.debug("[{}][{}] No activity events were received during reporting period for device with id: [{}]. Going to remove activity state.", - activityKey.getTenantId().getId(), name, activityKey.getDeviceId().getId()); - states.remove(activityKey); - } - if (activityState.getLastReportedTime() < lastRecordedTime) { - reporter.report(activityKey, lastRecordedTime, activityState, new ActivityReportCallback<>() { - @Override - public void onSuccess(IntegrationActivityKey key, long reportedTime) { - updateLastReportedTime(key, reportedTime); - } - - @Override - public void onFailure(IntegrationActivityKey key, Throwable t) { - log.debug("[{}][{}] Failed to report last activity event in a period for device with id: [{}].", - activityKey.getTenantId().getId(), name, activityKey.getDeviceId().getId()); - } - }); - } - } - } - - private void updateLastReportedTime(IntegrationActivityKey key, long newLastReportedTime) { - states.computeIfPresent(key, (__, activityState) -> { - activityState.setLastReportedTime(Math.max(activityState.getLastReportedTime(), newLastReportedTime)); - return activityState; - }); - } - -} diff --git a/application/src/main/resources/thingsboard.yml b/application/src/main/resources/thingsboard.yml index 57e55f4eda..2f710f3e92 100644 --- a/application/src/main/resources/thingsboard.yml +++ b/application/src/main/resources/thingsboard.yml @@ -871,12 +871,13 @@ transport: # Interval of periodic check for expired sessions and report of the changes to session last activity time report_timeout: "${TB_TRANSPORT_SESSIONS_REPORT_TIMEOUT:3000}" activity: - # This property specifies the strategy for reporting activity events within each reporting period. - # The accepted values are 'first', 'last', and 'first-and-last'. - # - 'first': Only the first activity event in each reporting period is reported. - # - 'last': Only the last activity event in the reporting period is reported. - # - 'first-and-last': Both the first and last activity events in the reporting period are reported. - reporting_strategy: "${TB_TRANSPORT_ACTIVITY_REPORTING_STRATEGY:last}" + # This property specifies the strategy for reporting activity events within each reporting period. + # The accepted values are 'FIRST', 'LAST', 'FIRST_AND_LAST' and 'ALL'. + # - 'FIRST': Only the first activity event in each reporting period is reported. + # - 'LAST': Only the last activity event in the reporting period is reported. + # - 'FIRST_AND_LAST': Both the first and last activity events in the reporting period are reported. + # - 'ALL': All activity events in the reporting period are reported. + reporting_strategy: "${TB_TRANSPORT_ACTIVITY_REPORTING_STRATEGY:LAST}" json: # Cast String data types to Numeric if possible when processing Telemetry/Attributes JSON type_cast_enabled: "${JSON_TYPE_CAST_ENABLED:true}" diff --git a/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/activity/AbstractActivityManager.java b/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/activity/AbstractActivityManager.java index dd8682aaed..065a0106f3 100644 --- a/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/activity/AbstractActivityManager.java +++ b/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/activity/AbstractActivityManager.java @@ -1,22 +1,22 @@ /** * 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. @@ -33,25 +33,31 @@ package org.thingsboard.server.common.transport.activity; import lombok.extern.slf4j.Slf4j; import org.thingsboard.common.util.ThingsBoardThreadFactory; import org.thingsboard.server.common.data.StringUtils; +import org.thingsboard.server.common.transport.activity.strategy.ActivityStrategy; +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.AtomicBoolean; import java.util.function.Supplier; @Slf4j -public abstract class AbstractActivityManager implements ActivityManager { +public abstract class AbstractActivityManager implements ActivityManager> { - protected String name; - protected long reportingPeriodMillis; - protected ActivityStateReporter reporter; + private final ConcurrentMap> states = new ConcurrentHashMap<>(); + private String name; + private long reportingPeriodMillis; + private ActivityStateReporter> reporter; private ScheduledExecutorService scheduler; private boolean initialized; @Override - public synchronized void init(String name, long reportingPeriodMillis, ActivityStateReporter reporter) { + public synchronized void init(String name, long reportingPeriodMillis, ActivityStateReporter> reporter) { if (!initialized) { this.name = StringUtils.notBlankOrDefault(name, "activity-manager"); log.info("[{}] initializing.", this.name); @@ -69,7 +75,7 @@ public abstract class AbstractActivityManager } @Override - public void onActivity(Key key, Supplier newStateSupplier) { + public void onActivity(Key key, Supplier> newStateSupplier) { if (!initialized) { log.error("[{}] Failed to process activity event: activity manager is not initialized.", name); return; @@ -86,14 +92,83 @@ public abstract class AbstractActivityManager doOnActivity(key, newStateSupplier); } - protected abstract void doOnActivity(Key key, Supplier newStateSupplier); + private boolean validate(Key key, Supplier> newStateSupplier) { + + } + + protected abstract ActivityStrategy getStrategy(); + + private void doOnActivity(Key key, Supplier> newStateSupplier) { + long newLastRecordedTime = System.currentTimeMillis(); + + var shouldReport = new AtomicBoolean(false); + var activityState = states.compute(key, (__, state) -> { + if (state == null) { + state = newStateSupplier.get(); + state.setStrategy(getStrategy()); + } + if (state.getLastRecordedTime() < newLastRecordedTime) { + state.setLastRecordedTime(newLastRecordedTime); + } + shouldReport.set(state.getStrategy().onActivity(state)); + return state; + }); + + if (shouldReport.get()) { + log.debug("[{}] Going to report first activity event for key: [{}].", name, key); + reporter.report(key, activityState.getLastRecordedTime(), activityState, new ActivityReportCallback<>() { + @Override + public void onSuccess(Key key, long reportedTime) { + updateLastReportedTime(key, reportedTime); + } + + @Override + public void onFailure(Key key, Throwable t) { + log.debug("[{}] Failed to report first activity event for key: [{}].", name, key, t); + } + }); + } + } private void onReportingPeriodEnd() { log.debug("[{}] Going to end reporting period.", name); - doOnReportingPeriodEnd(); + for (Map.Entry> entry : states.entrySet()) { + var key = entry.getKey(); + var state = entry.getValue(); + long lastRecordedTime = state.getLastRecordedTime(); + + boolean hasExpired = false; // TODO: implement state expiration + boolean shouldReport; + if (hasExpired) { + states.remove(key); + shouldReport = true; + } else { + shouldReport = state.getStrategy().onReportingPeriodEnd(state); + } + + if (shouldReport) { + log.debug("[{}] Going to report last activity event for key: [{}].", name, key); + reporter.report(key, lastRecordedTime, state, new ActivityReportCallback<>() { + @Override + public void onSuccess(Key key, long newLastReportedTime) { + updateLastReportedTime(key, newLastReportedTime); + } + + @Override + public void onFailure(Key key, Throwable t) { + log.debug("[{}] Failed to report last activity event in a period for key: [{}].", name, key, t); + } + }); + } + } } - protected abstract void doOnReportingPeriodEnd(); + private void updateLastReportedTime(Key key, long newLastReportedTime) { + states.computeIfPresent(key, (__, activityState) -> { + activityState.setLastReportedTime(Math.max(activityState.getLastReportedTime(), newLastReportedTime)); + return activityState; + }); + } @Override public synchronized void destroy() { diff --git a/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/activity/ActivityManager.java b/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/activity/ActivityManager.java index 8f006003ef..53228cb1a2 100644 --- a/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/activity/ActivityManager.java +++ b/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/activity/ActivityManager.java @@ -1,22 +1,22 @@ /** * 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. @@ -32,11 +32,11 @@ package org.thingsboard.server.common.transport.activity; import java.util.function.Supplier; -public interface ActivityManager { +public interface ActivityManager { - void init(String name, long reportingPeriodMillis, ActivityStateReporter reporter); + void init(String name, long reportingPeriodMillis, ActivityStateReporter> reporter); - void onActivity(Key key, Supplier newStateSupplier); + void onActivity(Key key, Supplier> newStateSupplier); void destroy(); 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 index 09285f9192..f0f607ddd3 100644 --- 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 @@ -31,11 +31,14 @@ package org.thingsboard.server.common.transport.activity; import lombok.Data; +import org.thingsboard.server.common.transport.activity.strategy.ActivityStrategy; @Data -public class ActivityState { +public class ActivityState { private volatile long lastRecordedTime; private volatile long lastReportedTime; + private volatile ActivityStrategy strategy; + private volatile Metadata metadata; } diff --git a/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/activity/strategy/ActivityStrategy.java b/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/activity/strategy/ActivityStrategy.java new file mode 100644 index 0000000000..58c8c6d165 --- /dev/null +++ b/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/activity/strategy/ActivityStrategy.java @@ -0,0 +1,11 @@ +package org.thingsboard.server.common.transport.activity.strategy; + +import org.thingsboard.server.common.transport.activity.ActivityState; + +public interface ActivityStrategy { + + boolean onActivity(ActivityState state); + + boolean onReportingPeriodEnd(ActivityState state); + +} diff --git a/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/activity/strategy/ActivityStrategyType.java b/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/activity/strategy/ActivityStrategyType.java new file mode 100644 index 0000000000..b351d4a59a --- /dev/null +++ b/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/activity/strategy/ActivityStrategyType.java @@ -0,0 +1,20 @@ +package org.thingsboard.server.common.transport.activity.strategy; + +public enum ActivityStrategyType { + + FIRST(new FirstEventActivityStrategy()), + LAST(new LastEventActivityStrategy()), + FIRST_AND_LAST(new FirstAndLastEventActivityStrategy()), + ALL(new AllEventsActivityStrategy()); + + private final ActivityStrategy strategy; + + ActivityStrategyType(ActivityStrategy strategy) { + this.strategy = strategy; + } + + public ActivityStrategy getStrategy() { + return strategy; + } + +} diff --git a/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/activity/strategy/AllEventsActivityStrategy.java b/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/activity/strategy/AllEventsActivityStrategy.java new file mode 100644 index 0000000000..add491c3fa --- /dev/null +++ b/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/activity/strategy/AllEventsActivityStrategy.java @@ -0,0 +1,17 @@ +package org.thingsboard.server.common.transport.activity.strategy; + +import org.thingsboard.server.common.transport.activity.ActivityState; + +public class AllEventsActivityStrategy implements ActivityStrategy { + + @Override + public boolean onActivity(ActivityState state) { + return state.getLastReportedTime() < state.getLastReportedTime(); + } + + @Override + public boolean onReportingPeriodEnd(ActivityState state) { + return state.getLastReportedTime() < state.getLastReportedTime(); + } + +} diff --git a/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/activity/strategy/FirstAndLastEventActivityStrategy.java b/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/activity/strategy/FirstAndLastEventActivityStrategy.java new file mode 100644 index 0000000000..ecfbaff005 --- /dev/null +++ b/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/activity/strategy/FirstAndLastEventActivityStrategy.java @@ -0,0 +1,25 @@ +package org.thingsboard.server.common.transport.activity.strategy; + +import org.thingsboard.server.common.transport.activity.ActivityState; + +public class FirstAndLastEventActivityStrategy implements ActivityStrategy { + + private volatile long firstEventTs; + + @Override + public boolean onActivity(ActivityState state) { + long lastRecordedTime = state.getLastRecordedTime(); + if (firstEventTs == 0L) { + firstEventTs = lastRecordedTime; + return state.getLastReportedTime() < firstEventTs; + } + return false; + } + + @Override + public boolean onReportingPeriodEnd(ActivityState state) { + firstEventTs = 0L; + return state.getLastReportedTime() < state.getLastRecordedTime(); + } + +} diff --git a/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/activity/strategy/FirstEventActivityStrategy.java b/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/activity/strategy/FirstEventActivityStrategy.java new file mode 100644 index 0000000000..2b787c3251 --- /dev/null +++ b/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/activity/strategy/FirstEventActivityStrategy.java @@ -0,0 +1,25 @@ +package org.thingsboard.server.common.transport.activity.strategy; + +import org.thingsboard.server.common.transport.activity.ActivityState; + +public class FirstEventActivityStrategy implements ActivityStrategy { + + private volatile long firstEventTs; + + @Override + public boolean onActivity(ActivityState state) { + long lastRecordedTime = state.getLastRecordedTime(); + if (firstEventTs == 0L) { + firstEventTs = lastRecordedTime; + return state.getLastReportedTime() < firstEventTs; + } + return false; + } + + @Override + public boolean onReportingPeriodEnd(ActivityState state) { + firstEventTs = 0L; + return false; + } + +} diff --git a/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/activity/strategy/LastEventActivityStrategy.java b/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/activity/strategy/LastEventActivityStrategy.java new file mode 100644 index 0000000000..57260ed5ca --- /dev/null +++ b/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/activity/strategy/LastEventActivityStrategy.java @@ -0,0 +1,17 @@ +package org.thingsboard.server.common.transport.activity.strategy; + +import org.thingsboard.server.common.transport.activity.ActivityState; + +public class LastEventActivityStrategy implements ActivityStrategy { + + @Override + public boolean onActivity(ActivityState state) { + return false; + } + + @Override + public boolean onReportingPeriodEnd(ActivityState state) { + return state.getLastReportedTime() < state.getLastRecordedTime(); + } + +}