diff --git a/application/src/main/java/org/thingsboard/server/service/integration/activity/AllIntegrationActivityManager.java b/application/src/main/java/org/thingsboard/server/service/integration/activity/AllEventsIntegrationActivityManager.java similarity index 98% rename from application/src/main/java/org/thingsboard/server/service/integration/activity/AllIntegrationActivityManager.java rename to application/src/main/java/org/thingsboard/server/service/integration/activity/AllEventsIntegrationActivityManager.java index 5cb8d1ea38..f2edb20a7f 100644 --- a/application/src/main/java/org/thingsboard/server/service/integration/activity/AllIntegrationActivityManager.java +++ b/application/src/main/java/org/thingsboard/server/service/integration/activity/AllEventsIntegrationActivityManager.java @@ -53,7 +53,7 @@ import java.util.function.Supplier; @Component @TbCoreComponent @ConditionalOnProperty(prefix = "integrations.activity", value = "reporting_strategy", havingValue = "all") -public class AllIntegrationActivityManager extends AbstractActivityManager { +public class AllEventsIntegrationActivityManager extends AbstractActivityManager { private final ConcurrentMap states = new ConcurrentHashMap<>(); 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/FirstAndLastEventsIntegrationActivityManager.java similarity index 98% rename from application/src/main/java/org/thingsboard/server/service/integration/activity/FirstAndLastIntegrationActivityManager.java rename to application/src/main/java/org/thingsboard/server/service/integration/activity/FirstAndLastEventsIntegrationActivityManager.java index 785e6f4155..1c2dba809e 100644 --- a/application/src/main/java/org/thingsboard/server/service/integration/activity/FirstAndLastIntegrationActivityManager.java +++ b/application/src/main/java/org/thingsboard/server/service/integration/activity/FirstAndLastEventsIntegrationActivityManager.java @@ -54,7 +54,7 @@ import java.util.function.Supplier; @Component @TbCoreComponent @ConditionalOnProperty(prefix = "integrations.activity", value = "reporting_strategy", havingValue = "first-and-last") -public class FirstAndLastIntegrationActivityManager extends AbstractActivityManager { +public class FirstAndLastEventsIntegrationActivityManager extends AbstractActivityManager { private final ConcurrentMap states = new ConcurrentHashMap<>(); 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/FirstEventOnlyIntegrationActivityManager.java similarity index 98% rename from application/src/main/java/org/thingsboard/server/service/integration/activity/FirstOnlyIntegrationActivityManager.java rename to application/src/main/java/org/thingsboard/server/service/integration/activity/FirstEventOnlyIntegrationActivityManager.java index d99745af0a..24f9a8f388 100644 --- a/application/src/main/java/org/thingsboard/server/service/integration/activity/FirstOnlyIntegrationActivityManager.java +++ b/application/src/main/java/org/thingsboard/server/service/integration/activity/FirstEventOnlyIntegrationActivityManager.java @@ -54,7 +54,7 @@ import java.util.function.Supplier; @Component @TbCoreComponent @ConditionalOnProperty(prefix = "integrations.activity", value = "reporting_strategy", havingValue = "first") -public class FirstOnlyIntegrationActivityManager extends AbstractActivityManager { +public class FirstEventOnlyIntegrationActivityManager extends AbstractActivityManager { private final ConcurrentMap states = new ConcurrentHashMap<>(); 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/LastEventOnlyIntegrationActivityManager.java similarity index 97% rename from application/src/main/java/org/thingsboard/server/service/integration/activity/LastOnlyIntegrationActivityManager.java rename to application/src/main/java/org/thingsboard/server/service/integration/activity/LastEventOnlyIntegrationActivityManager.java index c90b4f86e5..16699b2e3c 100644 --- a/application/src/main/java/org/thingsboard/server/service/integration/activity/LastOnlyIntegrationActivityManager.java +++ b/application/src/main/java/org/thingsboard/server/service/integration/activity/LastEventOnlyIntegrationActivityManager.java @@ -47,7 +47,7 @@ import java.util.function.Supplier; @Component @TbCoreComponent @ConditionalOnProperty(prefix = "integrations.activity", value = "reporting_strategy", havingValue = "last") -public class LastOnlyIntegrationActivityManager extends AbstractActivityManager { +public class LastEventOnlyIntegrationActivityManager extends AbstractActivityManager { private final ConcurrentMap states = new ConcurrentHashMap<>(); diff --git a/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/activity/AbstractActivityManager.java b/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/activity/AbstractActivityManager.java index 9d7b5c4e19..dd8682aaed 100644 --- a/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/activity/AbstractActivityManager.java +++ b/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/activity/AbstractActivityManager.java @@ -53,8 +53,8 @@ public abstract class AbstractActivityManager @Override public synchronized void init(String name, long reportingPeriodMillis, ActivityStateReporter reporter) { if (!initialized) { - log.info("[{}] initializing.", this.name); this.name = StringUtils.notBlankOrDefault(name, "activity-manager"); + log.info("[{}] initializing.", this.name); this.reporter = Objects.requireNonNull(reporter, "Failed to initialize activity manager: provided activity reporter is null."); if (reportingPeriodMillis <= 0) { reportingPeriodMillis = 3000; diff --git a/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/service/activity/FirstAndLastTransportActivityManager.java b/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/service/activity/FirstAndLastEventTransportActivityManager.java similarity index 94% rename from common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/service/activity/FirstAndLastTransportActivityManager.java rename to common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/service/activity/FirstAndLastEventTransportActivityManager.java index 2cc17a532e..9c582b2587 100644 --- a/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/service/activity/FirstAndLastTransportActivityManager.java +++ b/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/service/activity/FirstAndLastEventTransportActivityManager.java @@ -50,7 +50,9 @@ import org.thingsboard.server.common.transport.service.TransportActivityState; import org.thingsboard.server.gen.transport.TransportProtos; import org.thingsboard.server.queue.util.TbTransportComponent; +import java.util.HashSet; import java.util.Map; +import java.util.Set; import java.util.UUID; import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.ConcurrentMap; @@ -63,7 +65,7 @@ import static org.thingsboard.server.common.transport.service.DefaultTransportSe @Component @TbTransportComponent @ConditionalOnProperty(prefix = "transport.activity", value = "reporting_strategy", havingValue = "first-and-last") -public class FirstAndLastTransportActivityManager extends AbstractActivityManager { +public class FirstAndLastEventTransportActivityManager extends AbstractActivityManager { private final ConcurrentMap states = new ConcurrentHashMap<>(); @@ -129,6 +131,7 @@ public class FirstAndLastTransportActivityManager extends AbstractActivityManage @Override protected void doOnReportingPeriodEnd() { + Set statesToRemove = new HashSet<>(); for (Map.Entry entry : states.entrySet()) { var sessionId = entry.getKey(); var activityStateWrapper = entry.getValue(); @@ -138,8 +141,8 @@ public class FirstAndLastTransportActivityManager extends AbstractActivityManage if (sessionMetaData != null) { activityState.setSessionInfoProto(sessionMetaData.getSessionInfo()); } else { - log.debug("[{}] Session with id: [{}] is not present. Removing it's activity state.", name, sessionId); - states.remove(sessionId); + log.debug("[{}] Session with id: [{}] is not present. Marking it's activity state for removal.", name, sessionId); + statesToRemove.add(sessionId); } long lastActivityTime = activityState.getLastRecordedTime(); @@ -160,9 +163,9 @@ public class FirstAndLastTransportActivityManager extends AbstractActivityManage long expirationTime = System.currentTimeMillis() - sessionInactivityTimeout; boolean hasExpired = sessionMetaData != null && lastActivityTime < expirationTime; if (hasExpired) { - log.debug("[{}] Session with id: [{}] has expired due to last activity time: [{}]. Removing it's activity state.", name, sessionId, lastActivityTime); - states.remove(sessionId); + log.debug("[{}] Session with id: [{}] has expired due to last activity time: [{}]. Marking it's activity state for removal.", name, sessionId, lastActivityTime); transportService.deregisterSession(sessionInfo); + statesToRemove.add(sessionId); transportService.process(sessionInfo, SESSION_EVENT_MSG_CLOSED, null); sessionMetaData.getListener().onRemoteSessionCloseCommand(sessionId, SESSION_EXPIRED_NOTIFICATION_PROTO); } @@ -182,6 +185,7 @@ public class FirstAndLastTransportActivityManager extends AbstractActivityManage } activityStateWrapper.setAlreadyBeenReported(false); } + statesToRemove.forEach(states::remove); } private void updateLastReportedTime(UUID key, long newLastReportedTime) { diff --git a/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/service/activity/FirstOnlyTransportActivityManager.java b/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/service/activity/OnlyFirstEventTransportActivityManager.java similarity index 94% rename from common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/service/activity/FirstOnlyTransportActivityManager.java rename to common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/service/activity/OnlyFirstEventTransportActivityManager.java index 82d6fcccc8..19a6314450 100644 --- a/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/service/activity/FirstOnlyTransportActivityManager.java +++ b/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/service/activity/OnlyFirstEventTransportActivityManager.java @@ -50,7 +50,9 @@ import org.thingsboard.server.common.transport.service.TransportActivityState; import org.thingsboard.server.gen.transport.TransportProtos; import org.thingsboard.server.queue.util.TbTransportComponent; +import java.util.HashSet; import java.util.Map; +import java.util.Set; import java.util.UUID; import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.ConcurrentMap; @@ -63,7 +65,7 @@ import static org.thingsboard.server.common.transport.service.DefaultTransportSe @Component @TbTransportComponent @ConditionalOnProperty(prefix = "transport.activity", value = "reporting_strategy", havingValue = "first") -public class FirstOnlyTransportActivityManager extends AbstractActivityManager { +public class OnlyFirstEventTransportActivityManager extends AbstractActivityManager { private final ConcurrentMap states = new ConcurrentHashMap<>(); @@ -129,6 +131,7 @@ public class FirstOnlyTransportActivityManager extends AbstractActivityManager statesToRemove = new HashSet<>(); for (Map.Entry entry : states.entrySet()) { var sessionId = entry.getKey(); var activityStateWrapper = entry.getValue(); @@ -138,8 +141,8 @@ public class FirstOnlyTransportActivityManager extends AbstractActivityManager { +public class OnlyLastEventTransportActivityManager extends AbstractActivityManager { private final ConcurrentMap states = new ConcurrentHashMap<>(); @@ -67,9 +69,9 @@ public class LastOnlyTransportActivityManager extends AbstractActivityManager newStateSupplier) { + protected void doOnActivity(UUID sessionId, Supplier newStateSupplier) { long newLastRecordedTime = System.currentTimeMillis(); - states.compute(activityKey, (__, activityState) -> { + states.compute(sessionId, (__, activityState) -> { if (activityState == null) { activityState = newStateSupplier.get(); } @@ -82,6 +84,7 @@ public class LastOnlyTransportActivityManager extends AbstractActivityManager statesToRemove = new HashSet<>(); for (Map.Entry entry : states.entrySet()) { var sessionId = entry.getKey(); var activityState = entry.getValue(); @@ -90,8 +93,8 @@ public class LastOnlyTransportActivityManager extends AbstractActivityManager