Browse Source

Change activity manager names, remove activity after iteration on period end

pull/9980/head
Dmytro Skarzhynets 3 years ago
parent
commit
cfdceb9ae0
  1. 2
      application/src/main/java/org/thingsboard/server/service/integration/activity/AllEventsIntegrationActivityManager.java
  2. 2
      application/src/main/java/org/thingsboard/server/service/integration/activity/FirstAndLastEventsIntegrationActivityManager.java
  3. 2
      application/src/main/java/org/thingsboard/server/service/integration/activity/FirstEventOnlyIntegrationActivityManager.java
  4. 2
      application/src/main/java/org/thingsboard/server/service/integration/activity/LastEventOnlyIntegrationActivityManager.java
  5. 2
      common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/activity/AbstractActivityManager.java
  6. 14
      common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/service/activity/FirstAndLastEventTransportActivityManager.java
  7. 14
      common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/service/activity/OnlyFirstEventTransportActivityManager.java
  8. 18
      common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/service/activity/OnlyLastEventTransportActivityManager.java

2
application/src/main/java/org/thingsboard/server/service/integration/activity/AllIntegrationActivityManager.java → 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<IntegrationActivityKey, ActivityState> {
public class AllEventsIntegrationActivityManager extends AbstractActivityManager<IntegrationActivityKey, ActivityState> {
private final ConcurrentMap<IntegrationActivityKey, ActivityState> states = new ConcurrentHashMap<>();

2
application/src/main/java/org/thingsboard/server/service/integration/activity/FirstAndLastIntegrationActivityManager.java → 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<IntegrationActivityKey, ActivityState> {
public class FirstAndLastEventsIntegrationActivityManager extends AbstractActivityManager<IntegrationActivityKey, ActivityState> {
private final ConcurrentMap<IntegrationActivityKey, ActivityStateWrapper> states = new ConcurrentHashMap<>();

2
application/src/main/java/org/thingsboard/server/service/integration/activity/FirstOnlyIntegrationActivityManager.java → 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<IntegrationActivityKey, ActivityState> {
public class FirstEventOnlyIntegrationActivityManager extends AbstractActivityManager<IntegrationActivityKey, ActivityState> {
private final ConcurrentMap<IntegrationActivityKey, ActivityStateWrapper> states = new ConcurrentHashMap<>();

2
application/src/main/java/org/thingsboard/server/service/integration/activity/LastOnlyIntegrationActivityManager.java → 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<IntegrationActivityKey, ActivityState> {
public class LastEventOnlyIntegrationActivityManager extends AbstractActivityManager<IntegrationActivityKey, ActivityState> {
private final ConcurrentMap<IntegrationActivityKey, ActivityState> states = new ConcurrentHashMap<>();

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

@ -53,8 +53,8 @@ public abstract class AbstractActivityManager<Key, State extends ActivityState>
@Override
public synchronized void init(String name, long reportingPeriodMillis, ActivityStateReporter<Key, State> 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;

14
common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/service/activity/FirstAndLastTransportActivityManager.java → 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<UUID, TransportActivityState> {
public class FirstAndLastEventTransportActivityManager extends AbstractActivityManager<UUID, TransportActivityState> {
private final ConcurrentMap<UUID, ActivityStateWrapper> states = new ConcurrentHashMap<>();
@ -129,6 +131,7 @@ public class FirstAndLastTransportActivityManager extends AbstractActivityManage
@Override
protected void doOnReportingPeriodEnd() {
Set<UUID> statesToRemove = new HashSet<>();
for (Map.Entry<UUID, ActivityStateWrapper> 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) {

14
common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/service/activity/FirstOnlyTransportActivityManager.java → 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<UUID, TransportActivityState> {
public class OnlyFirstEventTransportActivityManager extends AbstractActivityManager<UUID, TransportActivityState> {
private final ConcurrentMap<UUID, ActivityStateWrapper> states = new ConcurrentHashMap<>();
@ -129,6 +131,7 @@ public class FirstOnlyTransportActivityManager extends AbstractActivityManager<U
@Override
protected void doOnReportingPeriodEnd() {
Set<UUID> statesToRemove = new HashSet<>();
for (Map.Entry<UUID, ActivityStateWrapper> entry : states.entrySet()) {
var sessionId = entry.getKey();
var activityStateWrapper = entry.getValue();
@ -138,8 +141,8 @@ public class FirstOnlyTransportActivityManager extends AbstractActivityManager<U
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 FirstOnlyTransportActivityManager extends AbstractActivityManager<U
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);
}
@ -183,6 +186,7 @@ public class FirstOnlyTransportActivityManager extends AbstractActivityManager<U
}
activityStateWrapper.setAlreadyBeenReported(false);
}
statesToRemove.forEach(states::remove);
}
private void updateLastReportedTime(UUID key, long newLastReportedTime) {

18
common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/service/activity/LastOnlyTransportActivityManager.java → common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/service/activity/OnlyLastEventTransportActivityManager.java

@ -43,7 +43,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;
@ -56,7 +58,7 @@ import static org.thingsboard.server.common.transport.service.DefaultTransportSe
@Component
@TbTransportComponent
@ConditionalOnProperty(prefix = "transport.activity", value = "reporting_strategy", havingValue = "last")
public class LastOnlyTransportActivityManager extends AbstractActivityManager<UUID, TransportActivityState> {
public class OnlyLastEventTransportActivityManager extends AbstractActivityManager<UUID, TransportActivityState> {
private final ConcurrentMap<UUID, TransportActivityState> states = new ConcurrentHashMap<>();
@ -67,9 +69,9 @@ public class LastOnlyTransportActivityManager extends AbstractActivityManager<UU
private TransportService transportService;
@Override
protected void doOnActivity(UUID activityKey, Supplier<TransportActivityState> newStateSupplier) {
protected void doOnActivity(UUID sessionId, Supplier<TransportActivityState> 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<UU
@Override
protected void doOnReportingPeriodEnd() {
Set<UUID> statesToRemove = new HashSet<>();
for (Map.Entry<UUID, TransportActivityState> entry : states.entrySet()) {
var sessionId = entry.getKey();
var activityState = entry.getValue();
@ -90,8 +93,8 @@ public class LastOnlyTransportActivityManager extends AbstractActivityManager<UU
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();
@ -112,9 +115,9 @@ public class LastOnlyTransportActivityManager extends AbstractActivityManager<UU
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);
}
@ -133,6 +136,7 @@ public class LastOnlyTransportActivityManager extends AbstractActivityManager<UU
});
}
}
statesToRemove.forEach(states::remove);
}
private void updateLastReportedTime(UUID key, long newLastReportedTime) {
Loading…
Cancel
Save