From 3461f6ae618a17529a8f0da39df16fac722716a1 Mon Sep 17 00:00:00 2001 From: Dmytro Skarzhynets Date: Thu, 7 Dec 2023 11:25:42 +0200 Subject: [PATCH] [WIP] Implemented three activity strategies for integration service --- ...irstAndLastIntegrationActivityManager.java | 196 ++++++++++++++ .../FirstOnlyIntegrationActivityManager.java | 192 ++++++++++++++ .../LastOnlyIntegrationActivityManager.java | 145 +++++++++++ ...StateManager.java => ActivityManager.java} | 4 +- .../activity/ActivityStateManagerImpl.java | 240 ------------------ .../activity/ActivityStateReportCallback.java | 2 - ...porter.java => ActivityStateReporter.java} | 6 +- .../service/DefaultTransportService.java | 46 ++-- .../LastOnlyTransportActivityManager.java | 131 ++++++++++ 9 files changed, 692 insertions(+), 270 deletions(-) create mode 100644 application/src/main/java/org/thingsboard/server/service/integration/activity/FirstAndLastIntegrationActivityManager.java create mode 100644 application/src/main/java/org/thingsboard/server/service/integration/activity/FirstOnlyIntegrationActivityManager.java create mode 100644 application/src/main/java/org/thingsboard/server/service/integration/activity/LastOnlyIntegrationActivityManager.java rename common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/activity/{ActivityStateManager.java => ActivityManager.java} (92%) delete mode 100644 common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/activity/ActivityStateManagerImpl.java rename common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/activity/{AsyncActivityStateReporter.java => ActivityStateReporter.java} (85%) create mode 100644 common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/service/activity/LastOnlyTransportActivityManager.java diff --git a/application/src/main/java/org/thingsboard/server/service/integration/activity/FirstAndLastIntegrationActivityManager.java b/application/src/main/java/org/thingsboard/server/service/integration/activity/FirstAndLastIntegrationActivityManager.java new file mode 100644 index 0000000000..12be2101c4 --- /dev/null +++ b/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 { + + private final ConcurrentMap activityStates = new ConcurrentHashMap<>(); + + @Data + private static class ActivityStateWrapper { + + volatile ActivityState activityState; + volatile boolean alreadyBeenReported; + + } + + private ScheduledExecutorService scheduler; + private final ActivityStateReporter reporter; + private final long reportingPeriodDurationMillis; + private final String name; + private boolean initialized; + + public FirstAndLastIntegrationActivityManager(ActivityStateReporter 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 newStateSupplier) { + long newLastRecordedTime = System.currentTimeMillis(); + SettableFuture> 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 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 statesToRemove = new HashSet<>(); + Set> entries = activityStates.entrySet(); + for (Map.Entry 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(); + } + } + } + } + +} diff --git a/application/src/main/java/org/thingsboard/server/service/integration/activity/FirstOnlyIntegrationActivityManager.java b/application/src/main/java/org/thingsboard/server/service/integration/activity/FirstOnlyIntegrationActivityManager.java new file mode 100644 index 0000000000..7a0bccd48c --- /dev/null +++ b/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 { + + private final ConcurrentMap activityStates = new ConcurrentHashMap<>(); + + @Data + private static class ActivityStateWrapper { + + volatile ActivityState activityState; + volatile boolean alreadyBeenReported; + + } + + private ScheduledExecutorService scheduler; + private final ActivityStateReporter reporter; + private final long reportingPeriodDurationMillis; + private final String name; + private boolean initialized; + + public FirstOnlyIntegrationActivityManager(ActivityStateReporter 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 newStateSupplier) { + long newLastRecordedTime = System.currentTimeMillis(); + SettableFuture> 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 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 statesToRemoveAndReport = new HashMap<>(); + Set> entries = activityStates.entrySet(); + for (Map.Entry 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(); + } + } + } + } + +} diff --git a/application/src/main/java/org/thingsboard/server/service/integration/activity/LastOnlyIntegrationActivityManager.java b/application/src/main/java/org/thingsboard/server/service/integration/activity/LastOnlyIntegrationActivityManager.java new file mode 100644 index 0000000000..0dd7aa6fc1 --- /dev/null +++ b/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 { + + private final ConcurrentMap activityStates = new ConcurrentHashMap<>(); + + private ScheduledExecutorService scheduler; + private final ActivityStateReporter reporter; + private final long reportingPeriodDurationMillis; + private final String name; + private boolean initialized; + + public LastOnlyIntegrationActivityManager(ActivityStateReporter 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 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 statesToRemove = new HashSet<>(); + Set> entries = activityStates.entrySet(); + for (Map.Entry 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(); + } + } + } + } + +} diff --git a/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/activity/ActivityStateManager.java b/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/activity/ActivityManager.java similarity index 92% rename from common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/activity/ActivityStateManager.java rename to common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/activity/ActivityManager.java index db24c2388d..d93756a922 100644 --- a/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/activity/ActivityStateManager.java +++ b/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/activity/ActivityManager.java @@ -32,11 +32,11 @@ package org.thingsboard.server.common.transport.activity; import java.util.function.Supplier; -public interface ActivityStateManager { +public interface ActivityManager { void init(); - void recordActivity(Key key, Supplier newActivityStateSupplier); + void onActivity(Key key, Supplier newStateSupplier); void destroy(); diff --git a/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/activity/ActivityStateManagerImpl.java b/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/activity/ActivityStateManagerImpl.java deleted file mode 100644 index 7d87b02f19..0000000000 --- a/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/activity/ActivityStateManagerImpl.java +++ /dev/null @@ -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 implements ActivityStateManager { - - private final ConcurrentMap activityStates = new ConcurrentHashMap<>(); - - @Data - private static class ActivityStateWrapper { - - volatile ActivityState activityState; - volatile boolean alreadyBeenReported; - - } - - private ScheduledExecutorService scheduler; - private final ReadWriteLock lock = new ReentrantReadWriteLock(); - private final AtomicInteger currentPeriodId = new AtomicInteger(); - private final AsyncActivityStateReporter reporter; - private final long reportingPeriodDurationMillis; - private final String name; - private boolean initialized; - - public ActivityStateManagerImpl(AsyncActivityStateReporter reporter, long reportingPeriodDurationMillis, String name) { - this.reporter = Objects.requireNonNull(reporter, "Failed to initialize activity manager: provided reporter is null."); - this.reportingPeriodDurationMillis = reportingPeriodDurationMillis; // TODO: add min/max duration validation - this.name = name == null ? "activity-state-manager" : name; - } - - @Override - public synchronized void init() { - if (!initialized) { - initialized = true; - scheduler = Executors.newSingleThreadScheduledExecutor(ThingsBoardThreadFactory.forName(name)); - scheduler.scheduleAtFixedRate(this::reportLastEventAndStartNewPeriod, new Random().nextInt((int) reportingPeriodDurationMillis), reportingPeriodDurationMillis, TimeUnit.MILLISECONDS); - } - } - - @Override - public void recordActivity(Key activityKey, Supplier newActivityStateSupplier) { - long newLastRecordedTime = System.currentTimeMillis(); - int capturedPeriodId = currentPeriodId.get(); - lock.readLock().lock(); - SettableFuture> reportCompletedFuture = SettableFuture.create(); - try { - activityStates.compute(activityKey, (key, activityStateWrapper) -> { - if (activityStateWrapper == null) { - State activityState = newActivityStateSupplier.get(); - activityState.setLastRecordedTime(newLastRecordedTime); - activityState.setLastReportedTime(0L); - activityStateWrapper = new ActivityStateWrapper(); - activityStateWrapper.setActivityState(activityState); - activityStateWrapper.setAlreadyBeenReported(false); - } else { - activityStateWrapper.getActivityState().setLastRecordedTime(newLastRecordedTime); - } - if (activityStateWrapper.isAlreadyBeenReported()) { - return activityStateWrapper; - } - var activityState = activityStateWrapper.getActivityState(); - if (activityState.getLastReportedTime() < activityState.getLastRecordedTime()) { - reporter.reportAsync(key, (State) activityState, new ActivityStateReportCallback<>() { - @Override - public void onSuccess(Key key, long reportedTime) { - reportCompletedFuture.set(Pair.of(key, reportedTime)); - } - - @Override - public void onFailure(Key key, Throwable t) { - reportCompletedFuture.setException(t); - } - - @Override - public void onRemove(Key key) { - lock.readLock().lock(); - try { - activityStates.remove(key); - } finally { - lock.readLock().unlock(); - } - } - }); - } - activityStateWrapper.setAlreadyBeenReported(true); - return activityStateWrapper; - }); - } finally { - lock.readLock().unlock(); - } - Futures.addCallback(reportCompletedFuture, new FutureCallback<>() { - @Override - public void onSuccess(Pair reportResult) { - lock.readLock().lock(); - try { - updateLastReportedTime(reportResult.getFirst(), reportResult.getSecond()); - } finally { - lock.readLock().unlock(); - } - } - - @Override - public void onFailure(@NonNull Throwable t) { // TODO: add failure logging - lock.readLock().lock(); - try { - rollbackReportedStatus(activityKey, capturedPeriodId, newLastRecordedTime); - } finally { - lock.readLock().unlock(); - } - } - }, MoreExecutors.directExecutor()); - } - - private void rollbackReportedStatus(Key key, int capturedPeriodId, long newLastRecordedTime) { - activityStates.computeIfPresent(key, (__, activityStateWrapper) -> { - if (capturedPeriodId == currentPeriodId.get() && activityStateWrapper.getActivityState().getLastReportedTime() < newLastRecordedTime) { - activityStateWrapper.setAlreadyBeenReported(false); - } - return activityStateWrapper; - }); - } - - private void reportLastEventAndStartNewPeriod() { - lock.writeLock().lock(); - try { - Map activityStateMap = activityStates.entrySet().stream() - .peek(entry -> entry.getValue().setAlreadyBeenReported(false)) - .collect(Collectors.toMap(Map.Entry::getKey, entry -> (State) entry.getValue().getActivityState())); - reporter.reportAsync(activityStateMap, new ActivityStateReportCallback<>() { - @Override - public void onSuccess(Key key, long reportedTime) { - lock.readLock().lock(); - try { - updateLastReportedTime(key, reportedTime); - } finally { - lock.readLock().unlock(); - } - } - - @Override - public void onFailure(Key key, Throwable t) { - log.debug("Failed to report activity state for key [{}]!", key, t); - } - - @Override - public void onRemove(Key key) { - lock.readLock().lock(); - try { - activityStates.remove(key); - } finally { - lock.readLock().unlock(); - } - } - }); - } finally { - currentPeriodId.incrementAndGet(); - lock.writeLock().unlock(); - } - } - - private void updateLastReportedTime(Key key, long newLastReportedTime) { - activityStates.computeIfPresent(key, (__, activityStateWrapper) -> { - var activityState = activityStateWrapper.getActivityState(); - activityState.setLastReportedTime(Math.max(activityState.getLastReportedTime(), newLastReportedTime)); - return activityStateWrapper; - }); - } - - @Override - public synchronized void destroy() { - if (initialized) { - initialized = false; - if (scheduler != null) { - scheduler.shutdown(); - try { - if (scheduler.awaitTermination(10L, TimeUnit.SECONDS)) { - scheduler.shutdownNow(); - } - } catch (InterruptedException e) { - scheduler.shutdownNow(); - Thread.currentThread().interrupt(); - } - } - } - } - -} diff --git a/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/activity/ActivityStateReportCallback.java b/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/activity/ActivityStateReportCallback.java index d9d897d4a1..4ff752d61d 100644 --- a/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/activity/ActivityStateReportCallback.java +++ b/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/activity/ActivityStateReportCallback.java @@ -36,6 +36,4 @@ public interface ActivityStateReportCallback { void onFailure(Key key, Throwable t); - void onRemove(Key key); - } diff --git a/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/activity/AsyncActivityStateReporter.java b/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/activity/ActivityStateReporter.java similarity index 85% rename from common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/activity/AsyncActivityStateReporter.java rename to common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/activity/ActivityStateReporter.java index ab9094c252..620646f894 100644 --- a/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/activity/AsyncActivityStateReporter.java +++ b/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/activity/ActivityStateReporter.java @@ -32,10 +32,10 @@ package org.thingsboard.server.common.transport.activity; import java.util.Map; -public interface AsyncActivityStateReporter { +public interface ActivityStateReporter { - void reportAsync(Key key, State activityState, ActivityStateReportCallback reportCallback); + void report(Key key, State activityState, ActivityStateReportCallback reportCallback); - void reportAsync(Map activityStates, ActivityStateReportCallback reportCallback); + void report(Map activityStates, ActivityStateReportCallback reportCallback); } diff --git a/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/service/DefaultTransportService.java b/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/service/DefaultTransportService.java index 7769762b84..31ab2df09e 100644 --- a/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/service/DefaultTransportService.java +++ b/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/service/DefaultTransportService.java @@ -72,16 +72,16 @@ import org.thingsboard.server.common.transport.TransportResourceCache; import org.thingsboard.server.common.transport.TransportService; import org.thingsboard.server.common.transport.TransportServiceCallback; import org.thingsboard.server.common.transport.TransportTenantProfileCache; -import org.thingsboard.server.common.transport.activity.ActivityStateManager; -import org.thingsboard.server.common.transport.activity.ActivityStateManagerImpl; +import org.thingsboard.server.common.transport.activity.ActivityManager; 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.TransportDeviceInfo; import org.thingsboard.server.common.transport.auth.ValidateDeviceCredentialsResponse; import org.thingsboard.server.common.transport.limits.EntityLimitKey; import org.thingsboard.server.common.transport.limits.EntityLimitsCache; 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.gen.transport.TransportProtos; import org.thingsboard.server.gen.transport.TransportProtos.ProvisionDeviceRequestMsg; @@ -203,7 +203,7 @@ public class DefaultTransportService implements TransportService { private ExecutorService mainConsumerExecutor; public final ConcurrentMap sessions = new ConcurrentHashMap<>(); - private ActivityStateManager activityStateManager; + private ActivityManager activityManager; private final Map toServerRpcPendingMap = new ConcurrentHashMap<>(); private volatile boolean stopped = false; @@ -253,8 +253,8 @@ public class DefaultTransportService implements TransportService { transportNotificationsConsumer.subscribe(Collections.singleton(tpi)); transportApiRequestTemplate.init(); mainConsumerExecutor = Executors.newSingleThreadExecutor(ThingsBoardThreadFactory.forName("transport-consumer")); - activityStateManager = new ActivityStateManagerImpl<>(reporter, sessionReportTimeout, "transport-activity-state-manager"); - activityStateManager.init(); + activityManager = new LastOnlyTransportActivityManager(reporter, sessionReportTimeout, "transport-activity-state-manager"); + activityManager.init(); } // 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 // (maybe having a queue of "first" events will help with this, queue will be cleared once event was successfully queued or later time was reported by event in next period) // setting alreadyBeenReported status in callbacks ensures that events are reported but can cause several "first" events to be reported - private final AsyncActivityStateReporter reporter = new AsyncActivityStateReporter<>() { + private final ActivityStateReporter reporter = new ActivityStateReporter<>() { @Override - public void reportAsync(UUID sessionId, TransportActivityState activityState, ActivityStateReportCallback reportCallback) { + public void report(UUID sessionId, TransportActivityState activityState, ActivityStateReportCallback reportCallback) { SessionMetaData sessionMetaData = sessions.get(sessionId); - TransportProtos.SubscriptionInfoProto subscriptionInfo = TransportProtos.SubscriptionInfoProto.newBuilder() - .setAttributeSubscription(sessionMetaData != null && sessionMetaData.isSubscribedToAttributes()) - .setRpcSubscription(sessionMetaData != null && sessionMetaData.isSubscribedToRPC()) - .setLastActivityTime(activityState.getLastRecordedTime()) - .build(); - process(activityState.getSessionInfoProto(), subscriptionInfo, new TransportServiceCallback<>() { + reportActivityStateToCore(activityState.getSessionInfoProto(), sessionMetaData, activityState.getLastRecordedTime(), new TransportServiceCallback<>() { @Override - public void onSuccess(Void msgAcknowledged) { + public void onSuccess(Void msg) { reportCallback.onSuccess(sessionId, activityState.getLastRecordedTime()); + } @Override public void onError(Throwable e) { reportCallback.onFailure(sessionId, e); + } }); } @Override - public void reportAsync(Map activityStates, ActivityStateReportCallback reportCallback) { + public void report(Map activityStates, ActivityStateReportCallback reportCallback) { long expirationTime = System.currentTimeMillis() - sessionInactivityTimeout; for (Map.Entry entry : activityStates.entrySet()) { var sessionId = entry.getKey(); var activityState = entry.getValue(); - long lastActivityTime = activityState.getLastRecordedTime(); - SessionMetaData sessionMetaData = sessions.get(sessionId); if (sessionMetaData != null) { activityState.setSessionInfoProto(sessionMetaData.getSessionInfo()); } else { - reportCallback.onRemove(sessionId); + log.info("Removing activity state due to session deregistration."); + activityStates.remove(sessionId); } + long lastActivityTime = activityState.getLastRecordedTime(); TransportProtos.SessionInfoProto sessionInfo = activityState.getSessionInfoProto(); 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); } 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); sessionMetaData.getListener().onRemoteSessionCloseCommand(sessionId, SESSION_EXPIRED_NOTIFICATION_PROTO); } else if (activityState.getLastReportedTime() < lastActivityTime) { @@ -406,8 +404,8 @@ public class DefaultTransportService implements TransportService { if (transportApiRequestTemplate != null) { transportApiRequestTemplate.stop(); } - if (activityStateManager != null) { - activityStateManager.destroy(); + if (activityManager != null) { + activityManager.destroy(); } } @@ -660,6 +658,7 @@ public class DefaultTransportService implements TransportService { @Override public void process(TransportProtos.SessionInfoProto sessionInfo, TransportProtos.SessionEventMsg msg, TransportServiceCallback callback) { + log.info("Received session event: [{}]", msg.getEvent()); if (checkLimits(sessionInfo, msg, callback)) { reportActivityInternal(sessionInfo); sendToDeviceActor(sessionInfo, TransportToDeviceActorMsg.newBuilder().setSessionInfo(sessionInfo) @@ -888,7 +887,7 @@ public class DefaultTransportService implements TransportService { } private void reportActivityInternal(TransportProtos.SessionInfoProto sessionInfo) { - activityStateManager.recordActivity(toSessionId(sessionInfo), () -> { + activityManager.onActivity(toSessionId(sessionInfo), () -> { var activityState = new TransportActivityState(); activityState.setSessionInfoProto(sessionInfo); return activityState; @@ -950,6 +949,7 @@ public class DefaultTransportService implements TransportService { log.debug("Stopping scheduler to avoid resending response if request has been ack."); currentSession.getScheduledFuture().cancel(false); } + log.info("Deregistering session."); sessions.remove(toSessionId(sessionInfo)); } diff --git a/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/service/activity/LastOnlyTransportActivityManager.java b/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/service/activity/LastOnlyTransportActivityManager.java new file mode 100644 index 0000000000..41f0207f72 --- /dev/null +++ b/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 { + + private final ConcurrentMap activityStates = new ConcurrentHashMap<>(); + + private ScheduledExecutorService scheduler; + private final ActivityStateReporter reporter; + private final long reportingPeriodDurationMillis; + private final String name; + private boolean initialized; + + public LastOnlyTransportActivityManager(ActivityStateReporter 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 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(); + } + } + } + } + +}