Browse Source

[WIP] Implemented three activity strategies for integration service

pull/9980/head
Dmytro Skarzhynets 3 years ago
parent
commit
3461f6ae61
  1. 196
      application/src/main/java/org/thingsboard/server/service/integration/activity/FirstAndLastIntegrationActivityManager.java
  2. 192
      application/src/main/java/org/thingsboard/server/service/integration/activity/FirstOnlyIntegrationActivityManager.java
  3. 145
      application/src/main/java/org/thingsboard/server/service/integration/activity/LastOnlyIntegrationActivityManager.java
  4. 4
      common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/activity/ActivityManager.java
  5. 240
      common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/activity/ActivityStateManagerImpl.java
  6. 2
      common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/activity/ActivityStateReportCallback.java
  7. 6
      common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/activity/ActivityStateReporter.java
  8. 46
      common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/service/DefaultTransportService.java
  9. 131
      common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/service/activity/LastOnlyTransportActivityManager.java

196
application/src/main/java/org/thingsboard/server/service/integration/activity/FirstAndLastIntegrationActivityManager.java

@ -0,0 +1,196 @@
/**
* 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.data.util.Pair;
import org.thingsboard.common.util.ThingsBoardThreadFactory;
import org.thingsboard.server.common.transport.activity.ActivityManager;
import org.thingsboard.server.common.transport.activity.ActivityState;
import org.thingsboard.server.common.transport.activity.ActivityStateReportCallback;
import org.thingsboard.server.common.transport.activity.ActivityStateReporter;
import java.util.HashSet;
import java.util.Map;
import java.util.Objects;
import java.util.Random;
import java.util.Set;
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.function.Supplier;
@Slf4j
public class FirstAndLastIntegrationActivityManager implements ActivityManager<IntegrationActivityKey, ActivityState> {
private final ConcurrentMap<IntegrationActivityKey, ActivityStateWrapper> activityStates = new ConcurrentHashMap<>();
@Data
private static class ActivityStateWrapper {
volatile ActivityState activityState;
volatile boolean alreadyBeenReported;
}
private ScheduledExecutorService scheduler;
private final ActivityStateReporter<IntegrationActivityKey, ActivityState> reporter;
private final long reportingPeriodDurationMillis;
private final String name;
private boolean initialized;
public FirstAndLastIntegrationActivityManager(ActivityStateReporter<IntegrationActivityKey, ActivityState> reporter, long reportingPeriodDurationMillis, String name) {
this.reporter = Objects.requireNonNull(reporter, "Failed to create 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) {
scheduler = Executors.newSingleThreadScheduledExecutor(ThingsBoardThreadFactory.forName(name));
scheduler.scheduleAtFixedRate(this::onReportingPeriodEnd, new Random().nextInt((int) reportingPeriodDurationMillis), reportingPeriodDurationMillis, TimeUnit.MILLISECONDS);
initialized = true;
}
}
@Override
public void onActivity(IntegrationActivityKey activityKey, Supplier<ActivityState> newStateSupplier) {
long newLastRecordedTime = System.currentTimeMillis();
SettableFuture<Pair<IntegrationActivityKey, Long>> reportCompletedFuture = SettableFuture.create();
activityStates.compute(activityKey, (key, activityStateWrapper) -> {
if (activityStateWrapper == null) {
ActivityState activityState = newStateSupplier.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.report(key, activityState, new ActivityStateReportCallback<>() {
@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) {
// TODO: log
}
}, MoreExecutors.directExecutor());
}
private void onReportingPeriodEnd() {
long expirationTime = System.currentTimeMillis() - reportingPeriodDurationMillis;
Set<IntegrationActivityKey> statesToRemove = new HashSet<>();
Set<Map.Entry<IntegrationActivityKey, ActivityStateWrapper>> entries = activityStates.entrySet();
for (Map.Entry<IntegrationActivityKey, ActivityStateWrapper> entry : entries) {
var activityKey = entry.getKey();
var activityStateWrapper = entry.getValue();
var activityState = activityStateWrapper.getActivityState();
if (activityState.getLastRecordedTime() < expirationTime) {
statesToRemove.add(activityKey);
}
if (activityState.getLastReportedTime() < activityState.getLastRecordedTime()) {
reporter.report(activityKey, activityState, new ActivityStateReportCallback<>() {
@Override
public void onSuccess(IntegrationActivityKey key, long reportedTime) {
updateLastReportedTime(key, reportedTime);
}
@Override
public void onFailure(IntegrationActivityKey key, Throwable t) {
// TODO: log
}
});
}
activityStateWrapper.setAlreadyBeenReported(false);
}
statesToRemove.forEach(activityStates::remove);
}
private void updateLastReportedTime(IntegrationActivityKey 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();
}
}
}
}
}

192
application/src/main/java/org/thingsboard/server/service/integration/activity/FirstOnlyIntegrationActivityManager.java

@ -0,0 +1,192 @@
/**
* 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.data.util.Pair;
import org.thingsboard.common.util.ThingsBoardThreadFactory;
import org.thingsboard.server.common.transport.activity.ActivityManager;
import org.thingsboard.server.common.transport.activity.ActivityState;
import org.thingsboard.server.common.transport.activity.ActivityStateReportCallback;
import org.thingsboard.server.common.transport.activity.ActivityStateReporter;
import java.util.HashMap;
import java.util.Map;
import java.util.Objects;
import java.util.Random;
import java.util.Set;
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.function.Supplier;
@Slf4j
public class FirstOnlyIntegrationActivityManager implements ActivityManager<IntegrationActivityKey, ActivityState> {
private final ConcurrentMap<IntegrationActivityKey, ActivityStateWrapper> activityStates = new ConcurrentHashMap<>();
@Data
private static class ActivityStateWrapper {
volatile ActivityState activityState;
volatile boolean alreadyBeenReported;
}
private ScheduledExecutorService scheduler;
private final ActivityStateReporter<IntegrationActivityKey, ActivityState> reporter;
private final long reportingPeriodDurationMillis;
private final String name;
private boolean initialized;
public FirstOnlyIntegrationActivityManager(ActivityStateReporter<IntegrationActivityKey, ActivityState> reporter, long reportingPeriodDurationMillis, String name) {
this.reporter = Objects.requireNonNull(reporter, "Failed to create 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) {
scheduler = Executors.newSingleThreadScheduledExecutor(ThingsBoardThreadFactory.forName(name));
scheduler.scheduleAtFixedRate(this::onReportingPeriodEnd, new Random().nextInt((int) reportingPeriodDurationMillis), reportingPeriodDurationMillis, TimeUnit.MILLISECONDS);
initialized = true;
}
}
@Override
public void onActivity(IntegrationActivityKey activityKey, Supplier<ActivityState> newStateSupplier) {
long newLastRecordedTime = System.currentTimeMillis();
SettableFuture<Pair<IntegrationActivityKey, Long>> reportCompletedFuture = SettableFuture.create();
activityStates.compute(activityKey, (key, activityStateWrapper) -> {
if (activityStateWrapper == null) {
ActivityState activityState = newStateSupplier.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.report(key, activityState, new ActivityStateReportCallback<>() {
@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) {
// TODO: log
}
}, MoreExecutors.directExecutor());
}
private void updateLastReportedTime(IntegrationActivityKey key, long newLastReportedTime) {
activityStates.computeIfPresent(key, (__, activityStateWrapper) -> {
var activityState = activityStateWrapper.getActivityState();
activityState.setLastReportedTime(Math.max(activityState.getLastReportedTime(), newLastReportedTime));
return activityStateWrapper;
});
}
private void onReportingPeriodEnd() {
Map<IntegrationActivityKey, ActivityState> statesToRemoveAndReport = new HashMap<>();
Set<Map.Entry<IntegrationActivityKey, ActivityStateWrapper>> entries = activityStates.entrySet();
for (Map.Entry<IntegrationActivityKey, ActivityStateWrapper> entry : entries) {
var activityKey = entry.getKey();
var activityStateWrapper = entry.getValue();
var activityState = activityStateWrapper.getActivityState();
if (!activityStateWrapper.isAlreadyBeenReported() && activityState.getLastReportedTime() < activityState.getLastRecordedTime()) {
statesToRemoveAndReport.put(activityKey, activityState);
}
activityStateWrapper.setAlreadyBeenReported(false);
}
statesToRemoveAndReport.forEach((key, state) -> reporter.report(key, state, new ActivityStateReportCallback<>() {
@Override
public void onSuccess(IntegrationActivityKey key, long reportedTime) {
activityStates.remove(key);
}
@Override
public void onFailure(IntegrationActivityKey key, Throwable t) {
// TODO: log
}
}));
}
@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();
}
}
}
}
}

145
application/src/main/java/org/thingsboard/server/service/integration/activity/LastOnlyIntegrationActivityManager.java

@ -0,0 +1,145 @@
/**
* 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.thingsboard.common.util.ThingsBoardThreadFactory;
import org.thingsboard.server.common.transport.activity.ActivityManager;
import org.thingsboard.server.common.transport.activity.ActivityState;
import org.thingsboard.server.common.transport.activity.ActivityStateReportCallback;
import org.thingsboard.server.common.transport.activity.ActivityStateReporter;
import java.util.HashSet;
import java.util.Map;
import java.util.Objects;
import java.util.Random;
import java.util.Set;
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.function.Supplier;
@Slf4j
public class LastOnlyIntegrationActivityManager implements ActivityManager<IntegrationActivityKey, ActivityState> {
private final ConcurrentMap<IntegrationActivityKey, ActivityState> activityStates = new ConcurrentHashMap<>();
private ScheduledExecutorService scheduler;
private final ActivityStateReporter<IntegrationActivityKey, ActivityState> reporter;
private final long reportingPeriodDurationMillis;
private final String name;
private boolean initialized;
public LastOnlyIntegrationActivityManager(ActivityStateReporter<IntegrationActivityKey, ActivityState> reporter, long reportingPeriodDurationMillis, String name) {
this.reporter = Objects.requireNonNull(reporter, "Failed to create 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) {
scheduler = Executors.newSingleThreadScheduledExecutor(ThingsBoardThreadFactory.forName(name));
scheduler.scheduleAtFixedRate(this::onReportingPeriodEnd, new Random().nextInt((int) reportingPeriodDurationMillis), reportingPeriodDurationMillis, TimeUnit.MILLISECONDS);
initialized = true;
}
}
@Override
public void onActivity(IntegrationActivityKey activityKey, Supplier<ActivityState> newStateSupplier) {
long newLastRecordedTime = System.currentTimeMillis();
activityStates.compute(activityKey, (key, activityState) -> {
if (activityState == null) {
activityState = newStateSupplier.get();
activityState.setLastRecordedTime(newLastRecordedTime);
activityState.setLastReportedTime(0L);
} else {
activityState.setLastRecordedTime(newLastRecordedTime);
}
return activityState;
});
}
private void onReportingPeriodEnd() {
long expirationTime = System.currentTimeMillis() - reportingPeriodDurationMillis;
Set<IntegrationActivityKey> statesToRemove = new HashSet<>();
Set<Map.Entry<IntegrationActivityKey, ActivityState>> entries = activityStates.entrySet();
for (Map.Entry<IntegrationActivityKey, ActivityState> entry : entries) {
var activityKey = entry.getKey();
var activityState = entry.getValue();
if (activityState.getLastRecordedTime() < expirationTime) {
statesToRemove.add(activityKey);
}
if (activityState.getLastReportedTime() < activityState.getLastRecordedTime()) {
reporter.report(activityKey, activityState, new ActivityStateReportCallback<>() {
@Override
public void onSuccess(IntegrationActivityKey key, long reportedTime) {
updateLastReportedTime(key, reportedTime);
}
@Override
public void onFailure(IntegrationActivityKey key, Throwable t) {
// TODO: log
}
});
}
}
statesToRemove.forEach(activityStates::remove);
}
private void updateLastReportedTime(IntegrationActivityKey key, long newLastReportedTime) {
activityStates.computeIfPresent(key, (__, activityState) -> {
activityState.setLastReportedTime(Math.max(activityState.getLastReportedTime(), newLastReportedTime));
return activityState;
});
}
@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();
}
}
}
}
}

4
common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/activity/ActivityStateManager.java → common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/activity/ActivityManager.java

@ -32,11 +32,11 @@ package org.thingsboard.server.common.transport.activity;
import java.util.function.Supplier; import java.util.function.Supplier;
public interface ActivityStateManager<Key, State extends ActivityState> { public interface ActivityManager<Key, State extends ActivityState> {
void init(); void init();
void recordActivity(Key key, Supplier<State> newActivityStateSupplier); void onActivity(Key key, Supplier<State> newStateSupplier);
void destroy(); void destroy();

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

@ -1,240 +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.common.transport.activity;
import com.google.common.util.concurrent.FutureCallback;
import com.google.common.util.concurrent.Futures;
import com.google.common.util.concurrent.MoreExecutors;
import com.google.common.util.concurrent.SettableFuture;
import lombok.Data;
import lombok.extern.slf4j.Slf4j;
import org.checkerframework.checker.nullness.qual.NonNull;
import org.springframework.data.util.Pair;
import org.thingsboard.common.util.ThingsBoardThreadFactory;
import java.util.Map;
import java.util.Objects;
import java.util.Random;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.ConcurrentMap;
import java.util.concurrent.Executors;
import java.util.concurrent.ScheduledExecutorService;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicInteger;
import java.util.concurrent.locks.ReadWriteLock;
import java.util.concurrent.locks.ReentrantReadWriteLock;
import java.util.function.Supplier;
import java.util.stream.Collectors;
@Slf4j
public final class ActivityStateManagerImpl<Key, State extends ActivityState> implements ActivityStateManager<Key, State> {
private final ConcurrentMap<Key, ActivityStateWrapper> activityStates = new ConcurrentHashMap<>();
@Data
private static class ActivityStateWrapper {
volatile ActivityState activityState;
volatile boolean alreadyBeenReported;
}
private ScheduledExecutorService scheduler;
private final ReadWriteLock lock = new ReentrantReadWriteLock();
private final AtomicInteger currentPeriodId = new AtomicInteger();
private final AsyncActivityStateReporter<Key, State> reporter;
private final long reportingPeriodDurationMillis;
private final String name;
private boolean initialized;
public ActivityStateManagerImpl(AsyncActivityStateReporter<Key, State> reporter, long reportingPeriodDurationMillis, String name) {
this.reporter = Objects.requireNonNull(reporter, "Failed to initialize activity manager: provided reporter is null.");
this.reportingPeriodDurationMillis = reportingPeriodDurationMillis; // TODO: add min/max duration validation
this.name = name == null ? "activity-state-manager" : name;
}
@Override
public synchronized void init() {
if (!initialized) {
initialized = true;
scheduler = Executors.newSingleThreadScheduledExecutor(ThingsBoardThreadFactory.forName(name));
scheduler.scheduleAtFixedRate(this::reportLastEventAndStartNewPeriod, new Random().nextInt((int) reportingPeriodDurationMillis), reportingPeriodDurationMillis, TimeUnit.MILLISECONDS);
}
}
@Override
public void recordActivity(Key activityKey, Supplier<State> newActivityStateSupplier) {
long newLastRecordedTime = System.currentTimeMillis();
int capturedPeriodId = currentPeriodId.get();
lock.readLock().lock();
SettableFuture<Pair<Key, Long>> reportCompletedFuture = SettableFuture.create();
try {
activityStates.compute(activityKey, (key, activityStateWrapper) -> {
if (activityStateWrapper == null) {
State activityState = newActivityStateSupplier.get();
activityState.setLastRecordedTime(newLastRecordedTime);
activityState.setLastReportedTime(0L);
activityStateWrapper = new ActivityStateWrapper();
activityStateWrapper.setActivityState(activityState);
activityStateWrapper.setAlreadyBeenReported(false);
} else {
activityStateWrapper.getActivityState().setLastRecordedTime(newLastRecordedTime);
}
if (activityStateWrapper.isAlreadyBeenReported()) {
return activityStateWrapper;
}
var activityState = activityStateWrapper.getActivityState();
if (activityState.getLastReportedTime() < activityState.getLastRecordedTime()) {
reporter.reportAsync(key, (State) activityState, new ActivityStateReportCallback<>() {
@Override
public void onSuccess(Key key, long reportedTime) {
reportCompletedFuture.set(Pair.of(key, reportedTime));
}
@Override
public void onFailure(Key key, Throwable t) {
reportCompletedFuture.setException(t);
}
@Override
public void onRemove(Key key) {
lock.readLock().lock();
try {
activityStates.remove(key);
} finally {
lock.readLock().unlock();
}
}
});
}
activityStateWrapper.setAlreadyBeenReported(true);
return activityStateWrapper;
});
} finally {
lock.readLock().unlock();
}
Futures.addCallback(reportCompletedFuture, new FutureCallback<>() {
@Override
public void onSuccess(Pair<Key, Long> reportResult) {
lock.readLock().lock();
try {
updateLastReportedTime(reportResult.getFirst(), reportResult.getSecond());
} finally {
lock.readLock().unlock();
}
}
@Override
public void onFailure(@NonNull Throwable t) { // TODO: add failure logging
lock.readLock().lock();
try {
rollbackReportedStatus(activityKey, capturedPeriodId, newLastRecordedTime);
} finally {
lock.readLock().unlock();
}
}
}, MoreExecutors.directExecutor());
}
private void rollbackReportedStatus(Key key, int capturedPeriodId, long newLastRecordedTime) {
activityStates.computeIfPresent(key, (__, activityStateWrapper) -> {
if (capturedPeriodId == currentPeriodId.get() && activityStateWrapper.getActivityState().getLastReportedTime() < newLastRecordedTime) {
activityStateWrapper.setAlreadyBeenReported(false);
}
return activityStateWrapper;
});
}
private void reportLastEventAndStartNewPeriod() {
lock.writeLock().lock();
try {
Map<Key, State> activityStateMap = activityStates.entrySet().stream()
.peek(entry -> entry.getValue().setAlreadyBeenReported(false))
.collect(Collectors.toMap(Map.Entry::getKey, entry -> (State) entry.getValue().getActivityState()));
reporter.reportAsync(activityStateMap, new ActivityStateReportCallback<>() {
@Override
public void onSuccess(Key key, long reportedTime) {
lock.readLock().lock();
try {
updateLastReportedTime(key, reportedTime);
} finally {
lock.readLock().unlock();
}
}
@Override
public void onFailure(Key key, Throwable t) {
log.debug("Failed to report activity state for key [{}]!", key, t);
}
@Override
public void onRemove(Key key) {
lock.readLock().lock();
try {
activityStates.remove(key);
} finally {
lock.readLock().unlock();
}
}
});
} finally {
currentPeriodId.incrementAndGet();
lock.writeLock().unlock();
}
}
private void updateLastReportedTime(Key key, long newLastReportedTime) {
activityStates.computeIfPresent(key, (__, activityStateWrapper) -> {
var activityState = activityStateWrapper.getActivityState();
activityState.setLastReportedTime(Math.max(activityState.getLastReportedTime(), newLastReportedTime));
return activityStateWrapper;
});
}
@Override
public synchronized void destroy() {
if (initialized) {
initialized = false;
if (scheduler != null) {
scheduler.shutdown();
try {
if (scheduler.awaitTermination(10L, TimeUnit.SECONDS)) {
scheduler.shutdownNow();
}
} catch (InterruptedException e) {
scheduler.shutdownNow();
Thread.currentThread().interrupt();
}
}
}
}
}

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

@ -36,6 +36,4 @@ public interface ActivityStateReportCallback<Key> {
void onFailure(Key key, Throwable t); void onFailure(Key key, Throwable t);
void onRemove(Key key);
} }

6
common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/activity/AsyncActivityStateReporter.java → common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/activity/ActivityStateReporter.java

@ -32,10 +32,10 @@ package org.thingsboard.server.common.transport.activity;
import java.util.Map; import java.util.Map;
public interface AsyncActivityStateReporter<Key, State extends ActivityState> { public interface ActivityStateReporter<Key, State extends ActivityState> {
void reportAsync(Key key, State activityState, ActivityStateReportCallback<Key> reportCallback); void report(Key key, State activityState, ActivityStateReportCallback<Key> reportCallback);
void reportAsync(Map<Key, State> activityStates, ActivityStateReportCallback<Key> reportCallback); void report(Map<Key, State> activityStates, ActivityStateReportCallback<Key> reportCallback);
} }

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

@ -72,16 +72,16 @@ import org.thingsboard.server.common.transport.TransportResourceCache;
import org.thingsboard.server.common.transport.TransportService; import org.thingsboard.server.common.transport.TransportService;
import org.thingsboard.server.common.transport.TransportServiceCallback; import org.thingsboard.server.common.transport.TransportServiceCallback;
import org.thingsboard.server.common.transport.TransportTenantProfileCache; import org.thingsboard.server.common.transport.TransportTenantProfileCache;
import org.thingsboard.server.common.transport.activity.ActivityStateManager; import org.thingsboard.server.common.transport.activity.ActivityManager;
import org.thingsboard.server.common.transport.activity.ActivityStateManagerImpl;
import org.thingsboard.server.common.transport.activity.ActivityStateReportCallback; import org.thingsboard.server.common.transport.activity.ActivityStateReportCallback;
import org.thingsboard.server.common.transport.activity.AsyncActivityStateReporter; import org.thingsboard.server.common.transport.activity.ActivityStateReporter;
import org.thingsboard.server.common.transport.auth.GetOrCreateDeviceFromGatewayResponse; import org.thingsboard.server.common.transport.auth.GetOrCreateDeviceFromGatewayResponse;
import org.thingsboard.server.common.transport.auth.TransportDeviceInfo; import org.thingsboard.server.common.transport.auth.TransportDeviceInfo;
import org.thingsboard.server.common.transport.auth.ValidateDeviceCredentialsResponse; import org.thingsboard.server.common.transport.auth.ValidateDeviceCredentialsResponse;
import org.thingsboard.server.common.transport.limits.EntityLimitKey; import org.thingsboard.server.common.transport.limits.EntityLimitKey;
import org.thingsboard.server.common.transport.limits.EntityLimitsCache; import org.thingsboard.server.common.transport.limits.EntityLimitsCache;
import org.thingsboard.server.common.transport.limits.TransportRateLimitService; import org.thingsboard.server.common.transport.limits.TransportRateLimitService;
import org.thingsboard.server.common.transport.service.activity.LastOnlyTransportActivityManager;
import org.thingsboard.server.common.transport.util.JsonUtils; import org.thingsboard.server.common.transport.util.JsonUtils;
import org.thingsboard.server.gen.transport.TransportProtos; import org.thingsboard.server.gen.transport.TransportProtos;
import org.thingsboard.server.gen.transport.TransportProtos.ProvisionDeviceRequestMsg; import org.thingsboard.server.gen.transport.TransportProtos.ProvisionDeviceRequestMsg;
@ -203,7 +203,7 @@ public class DefaultTransportService implements TransportService {
private ExecutorService mainConsumerExecutor; private ExecutorService mainConsumerExecutor;
public final ConcurrentMap<UUID, SessionMetaData> sessions = new ConcurrentHashMap<>(); public final ConcurrentMap<UUID, SessionMetaData> sessions = new ConcurrentHashMap<>();
private ActivityStateManager<UUID, TransportActivityState> activityStateManager; private ActivityManager<UUID, TransportActivityState> activityManager;
private final Map<String, RpcRequestMetadata> toServerRpcPendingMap = new ConcurrentHashMap<>(); private final Map<String, RpcRequestMetadata> toServerRpcPendingMap = new ConcurrentHashMap<>();
private volatile boolean stopped = false; private volatile boolean stopped = false;
@ -253,8 +253,8 @@ public class DefaultTransportService implements TransportService {
transportNotificationsConsumer.subscribe(Collections.singleton(tpi)); transportNotificationsConsumer.subscribe(Collections.singleton(tpi));
transportApiRequestTemplate.init(); transportApiRequestTemplate.init();
mainConsumerExecutor = Executors.newSingleThreadExecutor(ThingsBoardThreadFactory.forName("transport-consumer")); mainConsumerExecutor = Executors.newSingleThreadExecutor(ThingsBoardThreadFactory.forName("transport-consumer"));
activityStateManager = new ActivityStateManagerImpl<>(reporter, sessionReportTimeout, "transport-activity-state-manager"); activityManager = new LastOnlyTransportActivityManager(reporter, sessionReportTimeout, "transport-activity-state-manager");
activityStateManager.init(); activityManager.init();
} }
// TODO: why sessions and session activities are managed in a separate maps? maybe it is better to manage in a single map // TODO: why sessions and session activities are managed in a separate maps? maybe it is better to manage in a single map
@ -267,44 +267,41 @@ public class DefaultTransportService implements TransportService {
// this can cause lost activity updates if first event that got reported failed to enqueue in Kafka and next reporting period is already started // 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) // (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 // setting alreadyBeenReported status in callbacks ensures that events are reported but can cause several "first" events to be reported
private final AsyncActivityStateReporter<UUID, TransportActivityState> reporter = new AsyncActivityStateReporter<>() { private final ActivityStateReporter<UUID, TransportActivityState> reporter = new ActivityStateReporter<>() {
@Override @Override
public void reportAsync(UUID sessionId, TransportActivityState activityState, ActivityStateReportCallback<UUID> reportCallback) { public void report(UUID sessionId, TransportActivityState activityState, ActivityStateReportCallback<UUID> reportCallback) {
SessionMetaData sessionMetaData = sessions.get(sessionId); SessionMetaData sessionMetaData = sessions.get(sessionId);
TransportProtos.SubscriptionInfoProto subscriptionInfo = TransportProtos.SubscriptionInfoProto.newBuilder() reportActivityStateToCore(activityState.getSessionInfoProto(), sessionMetaData, activityState.getLastRecordedTime(), new TransportServiceCallback<>() {
.setAttributeSubscription(sessionMetaData != null && sessionMetaData.isSubscribedToAttributes())
.setRpcSubscription(sessionMetaData != null && sessionMetaData.isSubscribedToRPC())
.setLastActivityTime(activityState.getLastRecordedTime())
.build();
process(activityState.getSessionInfoProto(), subscriptionInfo, new TransportServiceCallback<>() {
@Override @Override
public void onSuccess(Void msgAcknowledged) { public void onSuccess(Void msg) {
reportCallback.onSuccess(sessionId, activityState.getLastRecordedTime()); reportCallback.onSuccess(sessionId, activityState.getLastRecordedTime());
} }
@Override @Override
public void onError(Throwable e) { public void onError(Throwable e) {
reportCallback.onFailure(sessionId, e); reportCallback.onFailure(sessionId, e);
} }
}); });
} }
@Override @Override
public void reportAsync(Map<UUID, TransportActivityState> activityStates, ActivityStateReportCallback<UUID> reportCallback) { public void report(Map<UUID, TransportActivityState> activityStates, ActivityStateReportCallback<UUID> reportCallback) {
long expirationTime = System.currentTimeMillis() - sessionInactivityTimeout; long expirationTime = System.currentTimeMillis() - sessionInactivityTimeout;
for (Map.Entry<UUID, TransportActivityState> entry : activityStates.entrySet()) { for (Map.Entry<UUID, TransportActivityState> entry : activityStates.entrySet()) {
var sessionId = entry.getKey(); var sessionId = entry.getKey();
var activityState = entry.getValue(); var activityState = entry.getValue();
long lastActivityTime = activityState.getLastRecordedTime();
SessionMetaData sessionMetaData = sessions.get(sessionId); SessionMetaData sessionMetaData = sessions.get(sessionId);
if (sessionMetaData != null) { if (sessionMetaData != null) {
activityState.setSessionInfoProto(sessionMetaData.getSessionInfo()); activityState.setSessionInfoProto(sessionMetaData.getSessionInfo());
} else { } else {
reportCallback.onRemove(sessionId); log.info("Removing activity state due to session deregistration.");
activityStates.remove(sessionId);
} }
long lastActivityTime = activityState.getLastRecordedTime();
TransportProtos.SessionInfoProto sessionInfo = activityState.getSessionInfoProto(); TransportProtos.SessionInfoProto sessionInfo = activityState.getSessionInfoProto();
if (sessionInfo.getGwSessionIdMSB() != 0 && sessionInfo.getGwSessionIdLSB() != 0) { if (sessionInfo.getGwSessionIdMSB() != 0 && sessionInfo.getGwSessionIdLSB() != 0) {
@ -323,7 +320,8 @@ public class DefaultTransportService implements TransportService {
log.debug("[{}] Session has expired due to last activity time: {}!", sessionId, lastActivityTime); log.debug("[{}] Session has expired due to last activity time: {}!", sessionId, lastActivityTime);
} }
sessions.remove(sessionId); sessions.remove(sessionId);
reportCallback.onRemove(sessionId); log.info("Removing activity state due to session expiration.");
activityStates.remove(sessionId);
process(sessionInfo, SESSION_EVENT_MSG_CLOSED, null); process(sessionInfo, SESSION_EVENT_MSG_CLOSED, null);
sessionMetaData.getListener().onRemoteSessionCloseCommand(sessionId, SESSION_EXPIRED_NOTIFICATION_PROTO); sessionMetaData.getListener().onRemoteSessionCloseCommand(sessionId, SESSION_EXPIRED_NOTIFICATION_PROTO);
} else if (activityState.getLastReportedTime() < lastActivityTime) { } else if (activityState.getLastReportedTime() < lastActivityTime) {
@ -406,8 +404,8 @@ public class DefaultTransportService implements TransportService {
if (transportApiRequestTemplate != null) { if (transportApiRequestTemplate != null) {
transportApiRequestTemplate.stop(); transportApiRequestTemplate.stop();
} }
if (activityStateManager != null) { if (activityManager != null) {
activityStateManager.destroy(); activityManager.destroy();
} }
} }
@ -660,6 +658,7 @@ public class DefaultTransportService implements TransportService {
@Override @Override
public void process(TransportProtos.SessionInfoProto sessionInfo, TransportProtos.SessionEventMsg msg, TransportServiceCallback<Void> callback) { public void process(TransportProtos.SessionInfoProto sessionInfo, TransportProtos.SessionEventMsg msg, TransportServiceCallback<Void> callback) {
log.info("Received session event: [{}]", msg.getEvent());
if (checkLimits(sessionInfo, msg, callback)) { if (checkLimits(sessionInfo, msg, callback)) {
reportActivityInternal(sessionInfo); reportActivityInternal(sessionInfo);
sendToDeviceActor(sessionInfo, TransportToDeviceActorMsg.newBuilder().setSessionInfo(sessionInfo) sendToDeviceActor(sessionInfo, TransportToDeviceActorMsg.newBuilder().setSessionInfo(sessionInfo)
@ -888,7 +887,7 @@ public class DefaultTransportService implements TransportService {
} }
private void reportActivityInternal(TransportProtos.SessionInfoProto sessionInfo) { private void reportActivityInternal(TransportProtos.SessionInfoProto sessionInfo) {
activityStateManager.recordActivity(toSessionId(sessionInfo), () -> { activityManager.onActivity(toSessionId(sessionInfo), () -> {
var activityState = new TransportActivityState(); var activityState = new TransportActivityState();
activityState.setSessionInfoProto(sessionInfo); activityState.setSessionInfoProto(sessionInfo);
return activityState; return activityState;
@ -950,6 +949,7 @@ public class DefaultTransportService implements TransportService {
log.debug("Stopping scheduler to avoid resending response if request has been ack."); log.debug("Stopping scheduler to avoid resending response if request has been ack.");
currentSession.getScheduledFuture().cancel(false); currentSession.getScheduledFuture().cancel(false);
} }
log.info("Deregistering session.");
sessions.remove(toSessionId(sessionInfo)); sessions.remove(toSessionId(sessionInfo));
} }

131
common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/service/activity/LastOnlyTransportActivityManager.java

@ -0,0 +1,131 @@
/**
* 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.activity;
import lombok.extern.slf4j.Slf4j;
import org.thingsboard.common.util.ThingsBoardThreadFactory;
import org.thingsboard.server.common.transport.activity.ActivityManager;
import org.thingsboard.server.common.transport.activity.ActivityStateReportCallback;
import org.thingsboard.server.common.transport.activity.ActivityStateReporter;
import org.thingsboard.server.common.transport.service.TransportActivityState;
import java.util.Objects;
import java.util.Random;
import java.util.UUID;
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.function.Supplier;
@Slf4j
public class LastOnlyTransportActivityManager implements ActivityManager<UUID, TransportActivityState> {
private final ConcurrentMap<UUID, TransportActivityState> activityStates = new ConcurrentHashMap<>();
private ScheduledExecutorService scheduler;
private final ActivityStateReporter<UUID, TransportActivityState> reporter;
private final long reportingPeriodDurationMillis;
private final String name;
private boolean initialized;
public LastOnlyTransportActivityManager(ActivityStateReporter<UUID, TransportActivityState> reporter, long reportingPeriodDurationMillis, String name) {
this.reporter = Objects.requireNonNull(reporter, "Failed to create 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) {
scheduler = Executors.newSingleThreadScheduledExecutor(ThingsBoardThreadFactory.forName(name));
scheduler.scheduleAtFixedRate(this::onReportingPeriodEnd, new Random().nextInt((int) reportingPeriodDurationMillis), reportingPeriodDurationMillis, TimeUnit.MILLISECONDS);
initialized = true;
}
}
@Override
public void onActivity(UUID activityKey, Supplier<TransportActivityState> newStateSupplier) {
long newLastRecordedTime = System.currentTimeMillis();
activityStates.compute(activityKey, (key, activityState) -> {
if (activityState == null) {
log.info("Creating new activity state.");
activityState = newStateSupplier.get();
activityState.setLastRecordedTime(newLastRecordedTime);
activityState.setLastReportedTime(0L);
} else {
activityState.setLastRecordedTime(newLastRecordedTime);
}
return activityState;
});
}
private void onReportingPeriodEnd() {
reporter.report(activityStates, new ActivityStateReportCallback<>() {
@Override
public void onSuccess(UUID key, long reportedTime) {
updateLastReportedTime(key, reportedTime);
}
@Override
public void onFailure(UUID uuid, Throwable t) {
// TODO: log
}
});
}
private void updateLastReportedTime(UUID key, long newLastReportedTime) {
activityStates.computeIfPresent(key, (__, activityState) -> {
activityState.setLastReportedTime(Math.max(activityState.getLastReportedTime(), newLastReportedTime));
return activityState;
});
}
@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();
}
}
}
}
}
Loading…
Cancel
Save