Browse Source

[WIP] Refactoring after review

pull/9980/head
Dmytro Skarzhynets 3 years ago
parent
commit
f93a49be6d
  1. 124
      application/src/main/java/org/thingsboard/server/service/integration/activity/AllEventsIntegrationActivityManager.java
  2. 157
      application/src/main/java/org/thingsboard/server/service/integration/activity/FirstAndLastEventsIntegrationActivityManager.java
  3. 160
      application/src/main/java/org/thingsboard/server/service/integration/activity/FirstEventOnlyIntegrationActivityManager.java
  4. 105
      application/src/main/java/org/thingsboard/server/service/integration/activity/LastEventOnlyIntegrationActivityManager.java
  5. 13
      application/src/main/resources/thingsboard.yml
  6. 103
      common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/activity/AbstractActivityManager.java
  7. 16
      common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/activity/ActivityManager.java
  8. 5
      common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/activity/ActivityState.java
  9. 11
      common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/activity/strategy/ActivityStrategy.java
  10. 20
      common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/activity/strategy/ActivityStrategyType.java
  11. 17
      common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/activity/strategy/AllEventsActivityStrategy.java
  12. 25
      common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/activity/strategy/FirstAndLastEventActivityStrategy.java
  13. 25
      common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/activity/strategy/FirstEventActivityStrategy.java
  14. 17
      common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/activity/strategy/LastEventActivityStrategy.java

124
application/src/main/java/org/thingsboard/server/service/integration/activity/AllEventsIntegrationActivityManager.java

@ -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<IntegrationActivityKey, ActivityState> {
private final ConcurrentMap<IntegrationActivityKey, ActivityState> states = new ConcurrentHashMap<>();
@Override
protected void doOnActivity(IntegrationActivityKey activityKey, Supplier<ActivityState> newStateSupplier) {
long newLastRecordedTime = System.currentTimeMillis();
SettableFuture<Pair<IntegrationActivityKey, Long>> 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<IntegrationActivityKey, Long> 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<IntegrationActivityKey, ActivityState> 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);
}
}
}
}

157
application/src/main/java/org/thingsboard/server/service/integration/activity/FirstAndLastEventsIntegrationActivityManager.java

@ -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<IntegrationActivityKey, ActivityState> {
private final ConcurrentMap<IntegrationActivityKey, ActivityStateWrapper> states = new ConcurrentHashMap<>();
@Data
private static class ActivityStateWrapper {
volatile ActivityState state;
volatile boolean alreadyBeenReported;
}
@Override
protected void doOnActivity(IntegrationActivityKey activityKey, Supplier<ActivityState> newStateSupplier) {
long newLastRecordedTime = System.currentTimeMillis();
SettableFuture<Pair<IntegrationActivityKey, Long>> 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<IntegrationActivityKey, Long> 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<IntegrationActivityKey, ActivityStateWrapper> 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;
});
}
}

160
application/src/main/java/org/thingsboard/server/service/integration/activity/FirstEventOnlyIntegrationActivityManager.java

@ -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<IntegrationActivityKey, ActivityState> {
private final ConcurrentMap<IntegrationActivityKey, ActivityStateWrapper> states = new ConcurrentHashMap<>();
@Data
private static class ActivityStateWrapper {
volatile ActivityState state;
volatile boolean alreadyBeenReported;
}
@Override
protected void doOnActivity(IntegrationActivityKey activityKey, Supplier<ActivityState> newStateSupplier) {
long newLastRecordedTime = System.currentTimeMillis();
SettableFuture<Pair<IntegrationActivityKey, Long>> 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<IntegrationActivityKey, Long> 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<IntegrationActivityKey, ActivityStateWrapper> 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;
});
}
}

105
application/src/main/java/org/thingsboard/server/service/integration/activity/LastEventOnlyIntegrationActivityManager.java

@ -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<IntegrationActivityKey, ActivityState> {
private final ConcurrentMap<IntegrationActivityKey, ActivityState> states = new ConcurrentHashMap<>();
@Override
protected void doOnActivity(IntegrationActivityKey activityKey, Supplier<ActivityState> 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<IntegrationActivityKey, ActivityState> 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;
});
}
}

13
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}"

103
common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/activity/AbstractActivityManager.java

@ -1,22 +1,22 @@
/**
* ThingsBoard, Inc. ("COMPANY") CONFIDENTIAL
*
* <p>
* Copyright © 2016-2023 ThingsBoard, Inc. All Rights Reserved.
*
* <p>
* 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.
*
* <p>
* Dissemination of this information or reproduction of this material is strictly forbidden
* unless prior written permission is obtained from COMPANY.
*
* <p>
* 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.
*
* <p>
* 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<Key, State extends ActivityState> implements ActivityManager<Key, State> {
public abstract class AbstractActivityManager<Key, Metadata> implements ActivityManager<Key, ActivityState<Metadata>> {
protected String name;
protected long reportingPeriodMillis;
protected ActivityStateReporter<Key, State> reporter;
private final ConcurrentMap<Key, ActivityState<Metadata>> states = new ConcurrentHashMap<>();
private String name;
private long reportingPeriodMillis;
private ActivityStateReporter<Key, ActivityState<Metadata>> reporter;
private ScheduledExecutorService scheduler;
private boolean initialized;
@Override
public synchronized void init(String name, long reportingPeriodMillis, ActivityStateReporter<Key, State> reporter) {
public synchronized void init(String name, long reportingPeriodMillis, ActivityStateReporter<Key, ActivityState<Metadata>> reporter) {
if (!initialized) {
this.name = StringUtils.notBlankOrDefault(name, "activity-manager");
log.info("[{}] initializing.", this.name);
@ -69,7 +75,7 @@ public abstract class AbstractActivityManager<Key, State extends ActivityState>
}
@Override
public void onActivity(Key key, Supplier<State> newStateSupplier) {
public void onActivity(Key key, Supplier<ActivityState<Metadata>> 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<Key, State extends ActivityState>
doOnActivity(key, newStateSupplier);
}
protected abstract void doOnActivity(Key key, Supplier<State> newStateSupplier);
private boolean validate(Key key, Supplier<ActivityState<Metadata>> newStateSupplier) {
}
protected abstract ActivityStrategy getStrategy();
private void doOnActivity(Key key, Supplier<ActivityState<Metadata>> 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<Key, ActivityState<Metadata>> 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() {

16
common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/activity/ActivityManager.java

@ -1,22 +1,22 @@
/**
* ThingsBoard, Inc. ("COMPANY") CONFIDENTIAL
*
* <p>
* Copyright © 2016-2023 ThingsBoard, Inc. All Rights Reserved.
*
* <p>
* 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.
*
* <p>
* Dissemination of this information or reproduction of this material is strictly forbidden
* unless prior written permission is obtained from COMPANY.
*
* <p>
* 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.
*
* <p>
* 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<Key, State extends ActivityState> {
public interface ActivityManager<Key, Metadata> {
void init(String name, long reportingPeriodMillis, ActivityStateReporter<Key, State> reporter);
void init(String name, long reportingPeriodMillis, ActivityStateReporter<Key, ActivityState<Metadata>> reporter);
void onActivity(Key key, Supplier<State> newStateSupplier);
void onActivity(Key key, Supplier<ActivityState<Metadata>> newStateSupplier);
void destroy();

5
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<Metadata> {
private volatile long lastRecordedTime;
private volatile long lastReportedTime;
private volatile ActivityStrategy strategy;
private volatile Metadata metadata;
}

11
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);
}

20
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;
}
}

17
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();
}
}

25
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();
}
}

25
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;
}
}

17
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();
}
}
Loading…
Cancel
Save