Browse Source

Merge pull request #9987 from thingsboard/feature/last-activity-time-granularity

Last activity time granularity
pull/9991/head
Andrew Shvayka 3 years ago
committed by GitHub
parent
commit
ace6ce9c28
No known key found for this signature in database GPG Key ID: 4AEE18F83AFDEB23
  1. 8
      application/src/main/java/org/thingsboard/server/service/state/DefaultDeviceStateService.java
  2. 14
      application/src/main/resources/thingsboard.yml
  3. 2
      common/transport/coap/src/main/java/org/thingsboard/server/transport/coap/client/DefaultCoapClientContext.java
  4. 10
      common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportHandler.java
  5. 2
      common/transport/snmp/src/main/java/org/thingsboard/server/transport/snmp/service/SnmpTransportService.java
  6. 2
      common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/TransportService.java
  7. 187
      common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/activity/AbstractActivityManager.java
  8. 26
      common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/activity/ActivityManager.java
  9. 24
      common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/activity/ActivityReportCallback.java
  10. 21
      common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/activity/ActivityState.java
  11. 24
      common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/activity/strategy/ActivityStrategy.java
  12. 47
      common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/activity/strategy/ActivityStrategyType.java
  13. 39
      common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/activity/strategy/AllEventsActivityStrategy.java
  14. 40
      common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/activity/strategy/FirstAndLastEventActivityStrategy.java
  15. 40
      common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/activity/strategy/FirstEventActivityStrategy.java
  16. 39
      common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/activity/strategy/LastEventActivityStrategy.java
  17. 126
      common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/service/DefaultTransportService.java
  18. 149
      common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/service/TransportActivityManager.java
  19. 44
      common/transport/transport-api/src/test/java/org/thingsboard/server/common/transport/activity/strategy/ActivityStrategyTypeTest.java
  20. 42
      common/transport/transport-api/src/test/java/org/thingsboard/server/common/transport/activity/strategy/AllEventsActivityStrategyTest.java
  21. 53
      common/transport/transport-api/src/test/java/org/thingsboard/server/common/transport/activity/strategy/FirstAndLastEventActivityStrategyTest.java
  22. 52
      common/transport/transport-api/src/test/java/org/thingsboard/server/common/transport/activity/strategy/FirstEventActivityStrategyTest.java
  23. 43
      common/transport/transport-api/src/test/java/org/thingsboard/server/common/transport/activity/strategy/LastEventActivityStrategyTest.java
  24. 508
      common/transport/transport-api/src/test/java/org/thingsboard/server/common/transport/service/TransportActivityManagerTest.java

8
application/src/main/java/org/thingsboard/server/service/state/DefaultDeviceStateService.java

@ -185,6 +185,10 @@ public class DefaultDeviceStateService extends AbstractPartitionBasedService<Dev
@Getter
private int initFetchPackSize;
@Value("${state.telemetryTtl:0}")
@Getter
private int telemetryTtl;
private ListeningExecutorService deviceStateExecutor;
final ConcurrentMap<DeviceId, DeviceStateData> deviceStates = new ConcurrentHashMap<>();
@ -803,7 +807,7 @@ public class DefaultDeviceStateService extends AbstractPartitionBasedService<Dev
tsSubService.saveAndNotifyInternal(
TenantId.SYS_TENANT_ID, deviceId,
Collections.singletonList(new BasicTsKvEntry(getCurrentTimeMillis(), new LongDataEntry(key, value))),
new TelemetrySaveCallback<>(deviceId, key, value));
telemetryTtl, new TelemetrySaveCallback<>(deviceId, key, value));
} else {
tsSubService.saveAttrAndNotify(TenantId.SYS_TENANT_ID, deviceId, SERVER_SCOPE, key, value, new TelemetrySaveCallback<>(deviceId, key, value));
}
@ -814,7 +818,7 @@ public class DefaultDeviceStateService extends AbstractPartitionBasedService<Dev
tsSubService.saveAndNotifyInternal(
TenantId.SYS_TENANT_ID, deviceId,
Collections.singletonList(new BasicTsKvEntry(getCurrentTimeMillis(), new BooleanDataEntry(key, value))),
new TelemetrySaveCallback<>(deviceId, key, value));
telemetryTtl, new TelemetrySaveCallback<>(deviceId, key, value));
} else {
tsSubService.saveAttrAndNotify(TenantId.SYS_TENANT_ID, deviceId, SERVER_SCOPE, key, value, new TelemetrySaveCallback<>(deviceId, key, value));
}

14
application/src/main/resources/thingsboard.yml

@ -785,6 +785,10 @@ state:
# If 'persistToTelemetry' is changed from 'false' to 'true': 'CREATE OR REPLACE VIEW device_info_view AS SELECT * FROM device_info_active_ts_view;'
# If 'persistToTelemetry' is changed from 'true' to 'false': 'CREATE OR REPLACE VIEW device_info_view AS SELECT * FROM device_info_active_attribute_view;'
persistToTelemetry: "${PERSIST_STATE_TO_TELEMETRY:false}"
# Millisecond value defining time-to-live for device state telemetry data (e.g. 'active', 'lastActivityTime').
# Used only when state.persistToTelemetry is set to 'true' and Cassandra is used for timeseries data.
# 0 means time-to-live mechanism is disabled.
telemetryTtl: "${STATE_TELEMETRY_TTL:0}"
# Tbel parameters
tbel:
@ -866,6 +870,14 @@ transport:
inactivity_timeout: "${TB_TRANSPORT_SESSIONS_INACTIVITY_TIMEOUT:300000}"
# Interval of periodic check for expired sessions and report of the changes to session last activity time
report_timeout: "${TB_TRANSPORT_SESSIONS_REPORT_TIMEOUT:3000}"
activity:
# This property specifies the strategy for reporting activity events within each reporting period.
# The accepted values are 'FIRST', 'LAST', 'FIRST_AND_LAST' and 'ALL'.
# - 'FIRST': Only the first activity event in each reporting period is reported.
# - 'LAST': Only the last activity event in the reporting period is reported.
# - 'FIRST_AND_LAST': Both the first and last activity events in the reporting period are reported.
# - 'ALL': All activity events in the reporting period are reported.
reporting_strategy: "${TB_TRANSPORT_ACTIVITY_REPORTING_STRATEGY:LAST}"
json:
# Cast String data types to Numeric if possible when processing Telemetry/Attributes JSON
type_cast_enabled: "${JSON_TYPE_CAST_ENABLED:true}"
@ -1695,4 +1707,4 @@ management:
web:
exposure:
# Expose metrics endpoint (use value 'prometheus' to enable prometheus metrics).
include: '${METRICS_ENDPOINTS_EXPOSE:info}'
include: '${METRICS_ENDPOINTS_EXPOSE:info}'

2
common/transport/coap/src/main/java/org/thingsboard/server/transport/coap/client/DefaultCoapClientContext.java

@ -212,7 +212,7 @@ public class DefaultCoapClientContext implements CoapClientContext {
public void reportActivity() {
for (TbCoapClientState state : clients.values()) {
if (state.getSession() != null) {
transportService.reportActivity(state.getSession());
transportService.recordActivity(state.getSession());
}
}
}

10
common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportHandler.java

@ -322,7 +322,7 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement
case PINGREQ:
if (checkConnected(ctx, msg)) {
ctx.writeAndFlush(new MqttMessage(new MqttFixedHeader(PINGRESP, false, AT_MOST_ONCE, false, 0)));
transportService.reportActivity(deviceSessionCtx.getSessionInfo());
transportService.recordActivity(deviceSessionCtx.getSessionInfo());
}
break;
case DISCONNECT:
@ -350,7 +350,7 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement
if (topicName.startsWith(MqttTopics.BASE_GATEWAY_API_TOPIC)) {
if (gatewaySessionHandler != null) {
handleGatewayPublishMsg(ctx, topicName, msgId, mqttMsg);
transportService.reportActivity(deviceSessionCtx.getSessionInfo());
transportService.recordActivity(deviceSessionCtx.getSessionInfo());
} else {
log.error("[gatewaySessionHandler] is null, [{}] Failed to process publish msg [{}][{}]", sessionId, topicName, msgId);
}
@ -526,7 +526,7 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement
transportService.process(deviceSessionCtx.getSessionInfo(), getAttributeMsg, getPubAckCallback(ctx, msgId, getAttributeMsg));
attrReqTopicType = TopicType.V2;
} else {
transportService.reportActivity(deviceSessionCtx.getSessionInfo());
transportService.recordActivity(deviceSessionCtx.getSessionInfo());
ack(ctx, msgId, ReturnCode.TOPIC_NAME_INVALID);
}
} catch (AdaptorException e) {
@ -796,7 +796,7 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement
}
}
if (!activityReported) {
transportService.reportActivity(deviceSessionCtx.getSessionInfo());
transportService.recordActivity(deviceSessionCtx.getSessionInfo());
}
ctx.writeAndFlush(createSubAckMessage(mqttMsg.variableHeader().messageId(), grantedQoSList));
}
@ -894,7 +894,7 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement
}
}
if (!activityReported) {
transportService.reportActivity(deviceSessionCtx.getSessionInfo());
transportService.recordActivity(deviceSessionCtx.getSessionInfo());
}
ctx.writeAndFlush(createUnSubAckMessage(mqttMsg.variableHeader().messageId(), unSubResults));
}

2
common/transport/snmp/src/main/java/org/thingsboard/server/transport/snmp/service/SnmpTransportService.java

@ -416,7 +416,7 @@ public class SnmpTransportService implements TbTransportService, CommandResponde
}
private void reportActivity(TransportProtos.SessionInfoProto sessionInfo) {
transportService.reportActivity(sessionInfo);
transportService.recordActivity(sessionInfo);
}

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

@ -145,7 +145,7 @@ public interface TransportService {
SessionMetaData registerSyncSession(SessionInfoProto sessionInfo, SessionMsgListener listener, long timeout);
void reportActivity(SessionInfoProto sessionInfo);
void recordActivity(SessionInfoProto sessionInfo);
void lifecycleEvent(TenantId tenantId, DeviceId deviceId, ComponentLifecycleEvent eventType, boolean success, Throwable error);

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

@ -0,0 +1,187 @@
/**
* Copyright © 2016-2024 The Thingsboard Authors
*
* Licensed under the Apache License, Version 2.0 (the "License");
* you may not use this file except in compliance with the License.
* You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package org.thingsboard.server.common.transport.activity;
import lombok.Data;
import lombok.extern.slf4j.Slf4j;
import org.springframework.beans.factory.annotation.Autowired;
import org.thingsboard.server.common.transport.activity.strategy.ActivityStrategy;
import org.thingsboard.server.queue.scheduler.SchedulerComponent;
import java.util.Map;
import java.util.Random;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.ConcurrentMap;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicBoolean;
import java.util.concurrent.atomic.AtomicLong;
import java.util.concurrent.atomic.AtomicReference;
@Slf4j
public abstract class AbstractActivityManager<Key, Metadata> implements ActivityManager<Key> {
private final ConcurrentMap<Key, ActivityStateWrapper> states = new ConcurrentHashMap<>();
@Autowired
protected SchedulerComponent scheduler;
@Data
private class ActivityStateWrapper {
private volatile ActivityState<Metadata> state;
private volatile long lastReportedTime;
private volatile ActivityStrategy strategy;
}
protected void init() {
var reportingPeriodMillis = getReportingPeriodMillis();
scheduler.scheduleAtFixedRate(this::onReportingPeriodEnd, new Random().nextInt((int) reportingPeriodMillis), reportingPeriodMillis, TimeUnit.MILLISECONDS);
}
protected abstract long getReportingPeriodMillis();
protected abstract ActivityState<Metadata> createNewState(Key key);
protected abstract ActivityStrategy getStrategy();
protected abstract ActivityState<Metadata> updateState(Key key, ActivityState<Metadata> state);
protected abstract boolean hasExpired(long lastRecordedTime);
protected abstract void onStateExpiry(Key key, Metadata metadata);
protected abstract void reportActivity(Key key, Metadata metadata, long timeToReport, ActivityReportCallback<Key> callback);
@Override
public void onActivity(Key key, long newLastRecordedTime) {
if (key == null) {
log.error("Failed to process activity event: provided activity key is null.");
return;
}
log.debug("Received activity event for key: [{}]", key);
var shouldReport = new AtomicBoolean(false);
var lastRecordedTime = new AtomicLong();
var lastReportedTime = new AtomicLong();
var metadata = new AtomicReference<Metadata>();
var activityStateWrapper = states.compute(key, (__, stateWrapper) -> {
if (stateWrapper == null) {
var newState = createNewState(key);
if (newState == null) {
return null;
}
stateWrapper = new ActivityStateWrapper();
stateWrapper.setState(newState);
stateWrapper.setStrategy(getStrategy());
}
var state = stateWrapper.getState();
if (state.getLastRecordedTime() < newLastRecordedTime) {
state.setLastRecordedTime(newLastRecordedTime);
}
shouldReport.set(stateWrapper.getStrategy().onActivity());
lastRecordedTime.set(state.getLastRecordedTime());
lastReportedTime.set(stateWrapper.getLastReportedTime());
metadata.set(state.getMetadata());
return stateWrapper;
});
if (activityStateWrapper == null) {
return;
}
if (shouldReport.get() && lastReportedTime.get() < lastRecordedTime.get()) {
log.debug("Going to report first activity event for key: [{}].", key);
reportActivity(key, metadata.get(), lastRecordedTime.get(), new ActivityReportCallback<>() {
@Override
public void onSuccess(Key key, long reportedTime) {
updateLastReportedTime(key, reportedTime);
}
@Override
public void onFailure(Key key, Throwable t) {
log.debug("Failed to report first activity event for key: [{}].", key, t);
}
});
}
}
@Override
public void onReportingPeriodEnd() {
log.debug("Going to end reporting period.");
for (Map.Entry<Key, ActivityStateWrapper> entry : states.entrySet()) {
var key = entry.getKey();
var stateWrapper = entry.getValue();
var currentState = stateWrapper.getState();
long lastRecordedTime = currentState.getLastRecordedTime();
long lastReportedTime = stateWrapper.getLastReportedTime();
var metadata = currentState.getMetadata();
boolean hasExpired;
boolean shouldReport;
var updatedState = updateState(key, currentState);
if (updatedState != null) {
stateWrapper.setState(updatedState);
lastRecordedTime = updatedState.getLastRecordedTime();
metadata = updatedState.getMetadata();
hasExpired = hasExpired(lastRecordedTime);
shouldReport = stateWrapper.getStrategy().onReportingPeriodEnd();
} else {
states.remove(key);
hasExpired = false;
shouldReport = true;
}
if (hasExpired) {
states.remove(key);
onStateExpiry(key, metadata);
shouldReport = true;
}
if (shouldReport && lastReportedTime < lastRecordedTime) {
log.debug("Going to report last activity event for key: [{}].", key);
reportActivity(key, metadata, lastRecordedTime, new ActivityReportCallback<>() {
@Override
public void onSuccess(Key key, long reportedTime) {
updateLastReportedTime(key, reportedTime);
}
@Override
public void onFailure(Key key, Throwable t) {
log.debug("Failed to report last activity event for key: [{}].", key, t);
}
});
}
}
}
@Override
public long getLastRecordedTime(Key key) {
ActivityStateWrapper stateWrapper = states.get(key);
return stateWrapper == null ? 0L : stateWrapper.getState().getLastRecordedTime();
}
private void updateLastReportedTime(Key key, long newLastReportedTime) {
states.computeIfPresent(key, (__, stateWrapper) -> {
stateWrapper.setLastReportedTime(Math.max(stateWrapper.getLastReportedTime(), newLastReportedTime));
return stateWrapper;
});
}
}

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

@ -0,0 +1,26 @@
/**
* Copyright © 2016-2024 The Thingsboard Authors
*
* Licensed under the Apache License, Version 2.0 (the "License");
* you may not use this file except in compliance with the License.
* You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package org.thingsboard.server.common.transport.activity;
public interface ActivityManager<Key> {
void onActivity(Key key, long activityTimeMillis);
void onReportingPeriodEnd();
long getLastRecordedTime(Key key);
}

24
common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/activity/ActivityReportCallback.java

@ -0,0 +1,24 @@
/**
* Copyright © 2016-2024 The Thingsboard Authors
*
* Licensed under the Apache License, Version 2.0 (the "License");
* you may not use this file except in compliance with the License.
* You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package org.thingsboard.server.common.transport.activity;
public interface ActivityReportCallback<Key> {
void onSuccess(Key key, long reportedTime);
void onFailure(Key key, Throwable t);
}

21
common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/service/SessionActivityData.java → common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/activity/ActivityState.java

@ -13,27 +13,14 @@
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package org.thingsboard.server.common.transport.service;
package org.thingsboard.server.common.transport.activity;
import lombok.Data;
import org.thingsboard.server.gen.transport.TransportProtos;
/**
* Created by ashvayka on 15.10.18.
*/
@Data
public class SessionActivityData {
private volatile TransportProtos.SessionInfoProto sessionInfo;
private volatile long lastActivityTime;
private volatile long lastReportedActivityTime;
SessionActivityData(TransportProtos.SessionInfoProto sessionInfo) {
this.sessionInfo = sessionInfo;
}
public class ActivityState<Metadata> {
void updateLastActivityTime() {
this.lastActivityTime = System.currentTimeMillis();
}
private volatile long lastRecordedTime;
private volatile Metadata metadata;
}

24
common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/activity/strategy/ActivityStrategy.java

@ -0,0 +1,24 @@
/**
* Copyright © 2016-2024 The Thingsboard Authors
*
* Licensed under the Apache License, Version 2.0 (the "License");
* you may not use this file except in compliance with the License.
* You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package org.thingsboard.server.common.transport.activity.strategy;
public interface ActivityStrategy {
boolean onActivity();
boolean onReportingPeriodEnd();
}

47
common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/activity/strategy/ActivityStrategyType.java

@ -0,0 +1,47 @@
/**
* Copyright © 2016-2024 The Thingsboard Authors
*
* Licensed under the Apache License, Version 2.0 (the "License");
* you may not use this file except in compliance with the License.
* You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package org.thingsboard.server.common.transport.activity.strategy;
public enum ActivityStrategyType {
ALL {
@Override
public ActivityStrategy toStrategy() {
return AllEventsActivityStrategy.getInstance();
}
},
FIRST {
@Override
public ActivityStrategy toStrategy() {
return new FirstEventActivityStrategy();
}
},
LAST {
@Override
public ActivityStrategy toStrategy() {
return LastEventActivityStrategy.getInstance();
}
},
FIRST_AND_LAST {
@Override
public ActivityStrategy toStrategy() {
return new FirstAndLastEventActivityStrategy();
}
};
public abstract ActivityStrategy toStrategy();
}

39
common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/activity/strategy/AllEventsActivityStrategy.java

@ -0,0 +1,39 @@
/**
* Copyright © 2016-2024 The Thingsboard Authors
*
* Licensed under the Apache License, Version 2.0 (the "License");
* you may not use this file except in compliance with the License.
* You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package org.thingsboard.server.common.transport.activity.strategy;
public final class AllEventsActivityStrategy implements ActivityStrategy {
private static final AllEventsActivityStrategy INSTANCE = new AllEventsActivityStrategy();
private AllEventsActivityStrategy() {
}
public static AllEventsActivityStrategy getInstance() {
return INSTANCE;
}
@Override
public boolean onActivity() {
return true;
}
@Override
public boolean onReportingPeriodEnd() {
return true;
}
}

40
common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/activity/strategy/FirstAndLastEventActivityStrategy.java

@ -0,0 +1,40 @@
/**
* Copyright © 2016-2024 The Thingsboard Authors
*
* Licensed under the Apache License, Version 2.0 (the "License");
* you may not use this file except in compliance with the License.
* You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package org.thingsboard.server.common.transport.activity.strategy;
import lombok.EqualsAndHashCode;
@EqualsAndHashCode
public final class FirstAndLastEventActivityStrategy implements ActivityStrategy {
private boolean firstEventReceived;
@Override
public synchronized boolean onActivity() {
if (!firstEventReceived) {
firstEventReceived = true;
return true;
}
return false;
}
@Override
public synchronized boolean onReportingPeriodEnd() {
firstEventReceived = false;
return true;
}
}

40
common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/activity/strategy/FirstEventActivityStrategy.java

@ -0,0 +1,40 @@
/**
* Copyright © 2016-2024 The Thingsboard Authors
*
* Licensed under the Apache License, Version 2.0 (the "License");
* you may not use this file except in compliance with the License.
* You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package org.thingsboard.server.common.transport.activity.strategy;
import lombok.EqualsAndHashCode;
@EqualsAndHashCode
public final class FirstEventActivityStrategy implements ActivityStrategy {
private boolean firstEventReceived;
@Override
public synchronized boolean onActivity() {
if (!firstEventReceived) {
firstEventReceived = true;
return true;
}
return false;
}
@Override
public synchronized boolean onReportingPeriodEnd() {
firstEventReceived = false;
return false;
}
}

39
common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/activity/strategy/LastEventActivityStrategy.java

@ -0,0 +1,39 @@
/**
* Copyright © 2016-2024 The Thingsboard Authors
*
* Licensed under the Apache License, Version 2.0 (the "License");
* you may not use this file except in compliance with the License.
* You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package org.thingsboard.server.common.transport.activity.strategy;
public final class LastEventActivityStrategy implements ActivityStrategy {
private static final LastEventActivityStrategy INSTANCE = new LastEventActivityStrategy();
private LastEventActivityStrategy() {
}
public static LastEventActivityStrategy getInstance() {
return INSTANCE;
}
@Override
public boolean onActivity() {
return false;
}
@Override
public boolean onReportingPeriodEnd() {
return true;
}
}

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

@ -50,9 +50,9 @@ import org.thingsboard.server.common.data.id.RuleChainId;
import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.common.data.id.TenantProfileId;
import org.thingsboard.server.common.data.limit.LimitedApi;
import org.thingsboard.server.common.data.msg.TbMsgType;
import org.thingsboard.server.common.data.notification.rule.trigger.RateLimitsTrigger;
import org.thingsboard.server.common.data.plugin.ComponentLifecycleEvent;
import org.thingsboard.server.common.data.msg.TbMsgType;
import org.thingsboard.server.common.data.rpc.RpcStatus;
import org.thingsboard.server.common.msg.TbMsg;
import org.thingsboard.server.common.msg.TbMsgMetaData;
@ -96,9 +96,9 @@ import org.thingsboard.server.queue.TbQueueProducer;
import org.thingsboard.server.queue.TbQueueRequestTemplate;
import org.thingsboard.server.queue.common.AsyncCallbackTemplate;
import org.thingsboard.server.queue.common.TbProtoQueueMsg;
import org.thingsboard.server.queue.discovery.TopicService;
import org.thingsboard.server.queue.discovery.PartitionService;
import org.thingsboard.server.queue.discovery.TbServiceInfoProvider;
import org.thingsboard.server.queue.discovery.TopicService;
import org.thingsboard.server.queue.provider.TbQueueProducerProvider;
import org.thingsboard.server.queue.provider.TbTransportQueueFactory;
import org.thingsboard.server.queue.scheduler.SchedulerComponent;
@ -109,16 +109,13 @@ import org.thingsboard.server.queue.util.TbTransportComponent;
import javax.annotation.PostConstruct;
import javax.annotation.PreDestroy;
import java.util.Collections;
import java.util.HashSet;
import java.util.LinkedHashMap;
import java.util.List;
import java.util.Map;
import java.util.Optional;
import java.util.Random;
import java.util.Set;
import java.util.UUID;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.ConcurrentMap;
import java.util.concurrent.ExecutionException;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;
@ -133,14 +130,12 @@ import java.util.stream.Collectors;
@Slf4j
@Service
@TbTransportComponent
public class DefaultTransportService implements TransportService {
public class DefaultTransportService extends TransportActivityManager implements TransportService {
public static final String OVERWRITE_ACTIVITY_TIME = "overwriteActivityTime";
public static final String SESSION_EXPIRED_MESSAGE = "Session has expired due to last activity time!";
public static final TransportProtos.SessionEventMsg SESSION_EVENT_MSG_OPEN = getSessionEventMsg(TransportProtos.SessionEvent.OPEN);
public static final TransportProtos.SessionEventMsg SESSION_EVENT_MSG_CLOSED = getSessionEventMsg(TransportProtos.SessionEvent.CLOSED);
public static final TransportProtos.SessionCloseNotificationProto SESSION_CLOSE_NOTIFICATION_PROTO = TransportProtos.SessionCloseNotificationProto.newBuilder()
.setMessage(SESSION_EXPIRED_MESSAGE).build();
public static final TransportProtos.SessionEventMsg SESSION_EVENT_MSG_OPEN = TransportProtos.SessionEventMsg.newBuilder()
.setSessionType(TransportProtos.SessionType.ASYNC)
.setEvent(TransportProtos.SessionEvent.OPEN).build();
public static final TransportProtos.SubscribeToAttributeUpdatesMsg SUBSCRIBE_TO_ATTRIBUTE_UPDATES_ASYNC_MSG = TransportProtos.SubscribeToAttributeUpdatesMsg.newBuilder()
.setSessionType(TransportProtos.SessionType.ASYNC).build();
public static final TransportProtos.SubscribeToRPCMsg SUBSCRIBE_TO_RPC_ASYNC_MSG = TransportProtos.SubscribeToRPCMsg.newBuilder()
@ -152,10 +147,6 @@ public class DefaultTransportService implements TransportService {
private boolean logEnabled;
@Value("${transport.log.max_length:1024}")
private int logMaxLength;
@Value("${transport.sessions.inactivity_timeout}")
private long sessionInactivityTimeout;
@Value("${transport.sessions.report_timeout}")
private long sessionReportTimeout;
@Value("${transport.client_side_rpc.timeout:60000}")
private long clientSideRpcTimeout;
@Value("${queue.transport.poll_interval}")
@ -182,11 +173,9 @@ public class DefaultTransportService implements TransportService {
private final TransportRateLimitService rateLimitService;
private final DataDecodingEncodingService dataDecodingEncodingService;
private final SchedulerComponent scheduler;
private final ApplicationEventPublisher eventPublisher;
private final TransportResourceCache transportResourceCache;
private final NotificationRuleProcessor notificationRuleProcessor;
private final EntityLimitsCache entityLimitsCache;
protected TbQueueRequestTemplate<TbProtoQueueMsg<TransportApiRequestMsg>, TbProtoQueueMsg<TransportApiResponseMsg>> transportApiRequestTemplate;
@ -201,8 +190,6 @@ public class DefaultTransportService implements TransportService {
protected ExecutorService transportCallbackExecutor;
private ExecutorService mainConsumerExecutor;
public final ConcurrentMap<UUID, SessionMetaData> sessions = new ConcurrentHashMap<>();
private final ConcurrentMap<UUID, SessionActivityData> sessionsActivity = new ConcurrentHashMap<>();
private final Map<String, RpcRequestMetadata> toServerRpcPendingMap = new ConcurrentHashMap<>();
private volatile boolean stopped = false;
@ -238,11 +225,11 @@ public class DefaultTransportService implements TransportService {
@PostConstruct
public void init() {
super.init();
this.ruleEngineProducerStats = statsFactory.createMessagesStats(StatsType.RULE_ENGINE.getName() + ".producer");
this.tbCoreProducerStats = statsFactory.createMessagesStats(StatsType.CORE.getName() + ".producer");
this.transportApiStats = statsFactory.createMessagesStats(StatsType.TRANSPORT.getName() + ".producer");
this.transportCallbackExecutor = ThingsBoardExecutors.newWorkStealingPool(20, getClass());
this.scheduler.scheduleAtFixedRate(this::checkInactivityAndReportActivity, new Random().nextInt((int) sessionReportTimeout), sessionReportTimeout, TimeUnit.MILLISECONDS);
this.scheduler.scheduleAtFixedRate(this::invalidateRateLimits, new Random().nextInt((int) sessionReportTimeout), sessionReportTimeout, TimeUnit.MILLISECONDS);
transportApiRequestTemplate = queueProvider.createTransportApiRequestTemplate();
transportApiRequestTemplate.setMessagesStats(transportApiStats);
@ -558,7 +545,7 @@ public class DefaultTransportService implements TransportService {
@Override
public void process(TransportProtos.SessionInfoProto sessionInfo, TransportProtos.SessionEventMsg msg, TransportServiceCallback<Void> callback) {
if (checkLimits(sessionInfo, msg, callback)) {
reportActivityInternal(sessionInfo);
recordActivityInternal(sessionInfo);
sendToDeviceActor(sessionInfo, TransportToDeviceActorMsg.newBuilder().setSessionInfo(sessionInfo)
.setSessionEvent(msg).build(), callback);
}
@ -578,7 +565,7 @@ public class DefaultTransportService implements TransportService {
}
}
reportActivityInternal(sessionInfo);
recordActivityInternal(sessionInfo);
sendToDeviceActor(sessionInfo, msg, callback);
}
}
@ -595,7 +582,7 @@ public class DefaultTransportService implements TransportService {
dataPoints += tsKv.getKvCount();
}
if (checkLimits(sessionInfo, msg, callback, dataPoints)) {
reportActivityInternal(sessionInfo);
recordActivityInternal(sessionInfo);
TenantId tenantId = getTenantId(sessionInfo);
DeviceId deviceId = new DeviceId(new UUID(sessionInfo.getDeviceIdMSB(), sessionInfo.getDeviceIdLSB()));
CustomerId customerId = getCustomerId(sessionInfo);
@ -619,7 +606,7 @@ public class DefaultTransportService implements TransportService {
@Override
public void process(TransportProtos.SessionInfoProto sessionInfo, TransportProtos.PostAttributeMsg msg, TbMsgMetaData md, TransportServiceCallback<Void> callback) {
if (checkLimits(sessionInfo, msg, callback, msg.getKvCount())) {
reportActivityInternal(sessionInfo);
recordActivityInternal(sessionInfo);
TenantId tenantId = getTenantId(sessionInfo);
DeviceId deviceId = new DeviceId(new UUID(sessionInfo.getDeviceIdMSB(), sessionInfo.getDeviceIdLSB()));
JsonObject json = JsonUtils.getJsonObject(msg.getKvList());
@ -639,7 +626,7 @@ public class DefaultTransportService implements TransportService {
@Override
public void process(TransportProtos.SessionInfoProto sessionInfo, TransportProtos.GetAttributeRequestMsg msg, TransportServiceCallback<Void> callback) {
if (checkLimits(sessionInfo, msg, callback)) {
reportActivityInternal(sessionInfo);
recordActivityInternal(sessionInfo);
sendToDeviceActor(sessionInfo, TransportToDeviceActorMsg.newBuilder().setSessionInfo(sessionInfo)
.setGetAttributes(msg).build(), new ApiStatsProxyCallback<>(getTenantId(sessionInfo), getCustomerId(sessionInfo), 1, callback));
}
@ -652,7 +639,7 @@ public class DefaultTransportService implements TransportService {
if (sessionMetaData != null) {
sessionMetaData.setSubscribedToAttributes(!msg.getUnsubscribe());
}
reportActivityInternal(sessionInfo);
recordActivityInternal(sessionInfo);
sendToDeviceActor(sessionInfo, TransportToDeviceActorMsg.newBuilder().setSessionInfo(sessionInfo).setSubscribeToAttributes(msg).build(),
new ApiStatsProxyCallback<>(getTenantId(sessionInfo), getCustomerId(sessionInfo), 1, callback));
}
@ -665,7 +652,7 @@ public class DefaultTransportService implements TransportService {
if (sessionMetaData != null) {
sessionMetaData.setSubscribedToRPC(!msg.getUnsubscribe());
}
reportActivityInternal(sessionInfo);
recordActivityInternal(sessionInfo);
sendToDeviceActor(sessionInfo, TransportToDeviceActorMsg.newBuilder().setSessionInfo(sessionInfo).setSubscribeToRPC(msg).build(),
new ApiStatsProxyCallback<>(getTenantId(sessionInfo), getCustomerId(sessionInfo), 1, callback));
}
@ -674,7 +661,7 @@ public class DefaultTransportService implements TransportService {
@Override
public void process(TransportProtos.SessionInfoProto sessionInfo, TransportProtos.ToDeviceRpcResponseMsg msg, TransportServiceCallback<Void> callback) {
if (checkLimits(sessionInfo, msg, callback)) {
reportActivityInternal(sessionInfo);
recordActivityInternal(sessionInfo);
sendToDeviceActor(sessionInfo, TransportToDeviceActorMsg.newBuilder().setSessionInfo(sessionInfo).setToDeviceRPCCallResponse(msg).build(),
new ApiStatsProxyCallback<>(getTenantId(sessionInfo), getCustomerId(sessionInfo), 1, callback));
}
@ -683,7 +670,7 @@ public class DefaultTransportService implements TransportService {
@Override
public void notifyAboutUplink(TransportProtos.SessionInfoProto sessionInfo, TransportProtos.UplinkNotificationMsg msg, TransportServiceCallback<Void> callback) {
if (checkLimits(sessionInfo, msg, callback)) {
reportActivityInternal(sessionInfo);
recordActivityInternal(sessionInfo);
sendToDeviceActor(sessionInfo, TransportToDeviceActorMsg.newBuilder().setSessionInfo(sessionInfo).setUplinkNotificationMsg(msg).build(), callback);
}
}
@ -704,7 +691,7 @@ public class DefaultTransportService implements TransportService {
if (checkLimits(sessionInfo, responseMsg, callback)) {
if (reportActivity) {
reportActivityInternal(sessionInfo);
recordActivityInternal(sessionInfo);
}
sendToDeviceActor(sessionInfo, TransportToDeviceActorMsg.newBuilder().setSessionInfo(sessionInfo).setRpcResponseStatusMsg(responseMsg).build(),
new ApiStatsProxyCallback<>(getTenantId(sessionInfo), getCustomerId(sessionInfo), 1, TransportServiceCallback.EMPTY));
@ -736,7 +723,7 @@ public class DefaultTransportService implements TransportService {
@Override
public void process(TransportProtos.SessionInfoProto sessionInfo, TransportProtos.ToServerRpcRequestMsg msg, TransportServiceCallback<Void> callback) {
if (checkLimits(sessionInfo, msg, callback)) {
reportActivityInternal(sessionInfo);
recordActivityInternal(sessionInfo);
UUID sessionId = toSessionId(sessionInfo);
TenantId tenantId = getTenantId(sessionInfo);
DeviceId deviceId = getDeviceId(sessionInfo);
@ -761,7 +748,7 @@ public class DefaultTransportService implements TransportService {
@Override
public void process(TransportProtos.SessionInfoProto sessionInfo, TransportProtos.ClaimDeviceMsg msg, TransportServiceCallback<Void> callback) {
if (checkLimits(sessionInfo, msg, callback)) {
reportActivityInternal(sessionInfo);
recordActivityInternal(sessionInfo);
sendToDeviceActor(sessionInfo, TransportToDeviceActorMsg.newBuilder().setSessionInfo(sessionInfo)
.setClaimDevice(msg).build(), callback);
}
@ -780,69 +767,16 @@ public class DefaultTransportService implements TransportService {
}
@Override
public void reportActivity(TransportProtos.SessionInfoProto sessionInfo) {
reportActivityInternal(sessionInfo);
public void recordActivity(TransportProtos.SessionInfoProto sessionInfo) {
recordActivityInternal(sessionInfo);
}
private void reportActivityInternal(TransportProtos.SessionInfoProto sessionInfo) {
UUID sessionId = toSessionId(sessionInfo);
SessionActivityData sessionMetaData = sessionsActivity.computeIfAbsent(sessionId, id -> new SessionActivityData(sessionInfo));
sessionMetaData.updateLastActivityTime();
}
private void checkInactivityAndReportActivity() {
long expTime = System.currentTimeMillis() - sessionInactivityTimeout;
Set<UUID> sessionsToRemove = new HashSet<>();
sessionsActivity.forEach((uuid, sessionAD) -> {
long lastActivityTime = sessionAD.getLastActivityTime();
SessionMetaData sessionMD = sessions.get(uuid);
if (sessionMD != null) {
sessionAD.setSessionInfo(sessionMD.getSessionInfo());
} else {
sessionsToRemove.add(uuid);
}
TransportProtos.SessionInfoProto sessionInfo = sessionAD.getSessionInfo();
if (sessionInfo.getGwSessionIdMSB() != 0 && sessionInfo.getGwSessionIdLSB() != 0) {
var gwSessionId = new UUID(sessionInfo.getGwSessionIdMSB(), sessionInfo.getGwSessionIdLSB());
SessionMetaData gwMetaData = sessions.get(gwSessionId);
SessionActivityData gwActivityData = sessionsActivity.get(gwSessionId);
if (gwMetaData != null && gwMetaData.isOverwriteActivityTime()) {
lastActivityTime = Math.max(gwActivityData.getLastActivityTime(), lastActivityTime);
}
}
if (lastActivityTime < expTime) {
if (sessionMD != null) {
if (log.isDebugEnabled()) {
log.debug("[{}] Session has expired due to last activity time: {}", toSessionId(sessionInfo), lastActivityTime);
}
sessions.remove(uuid);
sessionsToRemove.add(uuid);
process(sessionInfo, SESSION_EVENT_MSG_CLOSED, null);
sessionMD.getListener().onRemoteSessionCloseCommand(uuid, SESSION_CLOSE_NOTIFICATION_PROTO);
}
} else {
if (lastActivityTime > sessionAD.getLastReportedActivityTime()) {
final long lastActivityTimeFinal = lastActivityTime;
process(sessionInfo, TransportProtos.SubscriptionInfoProto.newBuilder()
.setAttributeSubscription(sessionMD != null && sessionMD.isSubscribedToAttributes())
.setRpcSubscription(sessionMD != null && sessionMD.isSubscribedToRPC())
.setLastActivityTime(lastActivityTime).build(), new TransportServiceCallback<Void>() {
@Override
public void onSuccess(Void msg) {
sessionAD.setLastReportedActivityTime(lastActivityTimeFinal);
}
private void recordActivityInternal(TransportProtos.SessionInfoProto sessionInfo) {
onActivity(toSessionId(sessionInfo), getCurrentTimeMillis());
}
@Override
public void onError(Throwable e) {
log.warn("[{}] Failed to report last activity time", uuid, e);
}
});
}
}
});
// Removes all closed or short-lived sessions.
sessionsToRemove.forEach(sessionsActivity::remove);
long getCurrentTimeMillis() {
return System.currentTimeMillis();
}
@Override
@ -913,7 +847,7 @@ public class DefaultTransportService implements TransportService {
}
TransportProtos.PostTelemetryMsg.Builder request = TransportProtos.PostTelemetryMsg.newBuilder();
TransportProtos.TsKvListProto.Builder builder = TransportProtos.TsKvListProto.newBuilder();
builder.setTs(TimeUnit.MILLISECONDS.toSeconds(System.currentTimeMillis()) * 1000L + (atomicTs.getAndIncrement() % 1000));
builder.setTs(TimeUnit.MILLISECONDS.toSeconds(getCurrentTimeMillis()) * 1000L + (atomicTs.getAndIncrement() % 1000));
builder.addKv(TransportProtos.KeyValueProto.newBuilder()
.setKey("transportLog")
.setType(TransportProtos.KeyValueType.STRING_V)
@ -1160,12 +1094,6 @@ public class DefaultTransportService implements TransportService {
return new DeviceId(new UUID(sessionInfo.getDeviceIdMSB(), sessionInfo.getDeviceIdLSB()));
}
private static TransportProtos.SessionEventMsg getSessionEventMsg(TransportProtos.SessionEvent event) {
return TransportProtos.SessionEventMsg.newBuilder()
.setSessionType(TransportProtos.SessionType.ASYNC)
.setEvent(event).build();
}
protected void sendToDeviceActor(TransportProtos.SessionInfoProto sessionInfo, TransportToDeviceActorMsg toDeviceActorMsg, TransportServiceCallback<Void> callback) {
ToCoreMsg toCoreMsg = ToCoreMsg.newBuilder().setToDeviceActorMsg(toDeviceActorMsg).build();
sendToCore(getTenantId(sessionInfo), getDeviceId(sessionInfo), toCoreMsg, getRoutingKey(sessionInfo), callback);

149
common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/service/TransportActivityManager.java

@ -0,0 +1,149 @@
/**
* Copyright © 2016-2024 The Thingsboard Authors
*
* Licensed under the Apache License, Version 2.0 (the "License");
* you may not use this file except in compliance with the License.
* You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package org.thingsboard.server.common.transport.service;
import lombok.extern.slf4j.Slf4j;
import org.springframework.beans.factory.annotation.Value;
import org.thingsboard.server.common.transport.TransportService;
import org.thingsboard.server.common.transport.TransportServiceCallback;
import org.thingsboard.server.common.transport.activity.AbstractActivityManager;
import org.thingsboard.server.common.transport.activity.ActivityReportCallback;
import org.thingsboard.server.common.transport.activity.ActivityState;
import org.thingsboard.server.common.transport.activity.strategy.ActivityStrategy;
import org.thingsboard.server.common.transport.activity.strategy.ActivityStrategyType;
import org.thingsboard.server.gen.transport.TransportProtos;
import java.util.UUID;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.ConcurrentMap;
@Slf4j
public abstract class TransportActivityManager extends AbstractActivityManager<UUID, TransportProtos.SessionInfoProto> implements TransportService {
public static final String SESSION_EXPIRED_MESSAGE = "Session has expired due to last activity time!";
public static final TransportProtos.SessionEventMsg SESSION_EVENT_MSG_CLOSED = TransportProtos.SessionEventMsg.newBuilder()
.setSessionType(TransportProtos.SessionType.ASYNC)
.setEvent(TransportProtos.SessionEvent.CLOSED).build();
public static final TransportProtos.SessionCloseNotificationProto SESSION_EXPIRED_NOTIFICATION_PROTO = TransportProtos.SessionCloseNotificationProto.newBuilder()
.setMessage(SESSION_EXPIRED_MESSAGE).build();
public final ConcurrentMap<UUID, SessionMetaData> sessions = new ConcurrentHashMap<>();
@Value("${transport.sessions.report_timeout}")
protected long sessionReportTimeout;
@Value("${transport.sessions.inactivity_timeout}")
protected long sessionInactivityTimeout;
@Value("${transport.activity.reporting_strategy:LAST}")
private ActivityStrategyType reportingStrategyType;
@Override
protected long getReportingPeriodMillis() {
return sessionReportTimeout;
}
@Override
protected ActivityState<TransportProtos.SessionInfoProto> createNewState(UUID sessionId) {
SessionMetaData session = sessions.get(sessionId);
if (session == null) {
return null;
}
ActivityState<TransportProtos.SessionInfoProto> state = new ActivityState<>();
state.setMetadata(session.getSessionInfo());
return state;
}
@Override
protected ActivityStrategy getStrategy() {
return reportingStrategyType.toStrategy();
}
@Override
protected ActivityState<TransportProtos.SessionInfoProto> updateState(UUID sessionId, ActivityState<TransportProtos.SessionInfoProto> state) {
SessionMetaData session = sessions.get(sessionId);
if (session == null) {
return null;
}
state.setMetadata(session.getSessionInfo());
var sessionInfo = state.getMetadata();
if (sessionInfo.getGwSessionIdMSB() == 0L || sessionInfo.getGwSessionIdLSB() == 0L) {
return state;
}
var gwSessionId = new UUID(sessionInfo.getGwSessionIdMSB(), sessionInfo.getGwSessionIdLSB());
SessionMetaData gwSession = sessions.get(gwSessionId);
if (gwSession == null || !gwSession.isOverwriteActivityTime()) {
return state;
}
long lastRecordedTime = state.getLastRecordedTime();
long gwLastRecordedTime = getLastRecordedTime(gwSessionId);
log.debug("Session with id: [{}] has gateway session with id: [{}] with overwrite activity time enabled. " +
"Updating last activity time. Session last recorded time: [{}], gateway session last recorded time: [{}].",
sessionId, gwSessionId, lastRecordedTime, gwLastRecordedTime);
state.setLastRecordedTime(Math.max(lastRecordedTime, gwLastRecordedTime));
return state;
}
@Override
protected boolean hasExpired(long lastRecordedTime) {
return (getCurrentTimeMillis() - sessionInactivityTimeout) > lastRecordedTime;
}
@Override
protected void onStateExpiry(UUID sessionId, TransportProtos.SessionInfoProto sessionInfo) {
log.debug("Session with id: [{}] has expired due to last activity time.", sessionId);
SessionMetaData expiredSession = sessions.remove(sessionId);
if (expiredSession != null) {
deregisterSession(sessionInfo);
process(sessionInfo, SESSION_EVENT_MSG_CLOSED, null);
expiredSession.getListener().onRemoteSessionCloseCommand(sessionId, SESSION_EXPIRED_NOTIFICATION_PROTO);
}
}
@Override
protected void reportActivity(UUID sessionId, TransportProtos.SessionInfoProto currentSessionInfo, long timeToReport, ActivityReportCallback<UUID> callback) {
log.debug("Reporting activity state for session with id: [{}]. Time to report: [{}].", sessionId, timeToReport);
SessionMetaData session = sessions.get(sessionId);
TransportProtos.SubscriptionInfoProto subscriptionInfo = TransportProtos.SubscriptionInfoProto.newBuilder()
.setAttributeSubscription(session != null && session.isSubscribedToAttributes())
.setRpcSubscription(session != null && session.isSubscribedToRPC())
.setLastActivityTime(timeToReport)
.build();
TransportProtos.SessionInfoProto sessionInfo = session != null ? session.getSessionInfo() : currentSessionInfo;
process(sessionInfo, subscriptionInfo, new TransportServiceCallback<>() {
@Override
public void onSuccess(Void msgAcknowledged) {
callback.onSuccess(sessionId, timeToReport);
}
@Override
public void onError(Throwable e) {
callback.onFailure(sessionId, e);
}
});
}
long getCurrentTimeMillis() {
return System.currentTimeMillis();
}
}

44
common/transport/transport-api/src/test/java/org/thingsboard/server/common/transport/activity/strategy/ActivityStrategyTypeTest.java

@ -0,0 +1,44 @@
/**
* Copyright © 2016-2024 The Thingsboard Authors
*
* Licensed under the Apache License, Version 2.0 (the "License");
* you may not use this file except in compliance with the License.
* You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package org.thingsboard.server.common.transport.activity.strategy;
import org.junit.jupiter.api.Test;
import static org.assertj.core.api.Assertions.assertThat;
public class ActivityStrategyTypeTest {
@Test
public void testCreateAllEventsStrategy() {
assertThat(ActivityStrategyType.ALL.toStrategy()).isEqualTo(AllEventsActivityStrategy.getInstance());
}
@Test
public void testCreateFirstEventStrategy() {
assertThat(ActivityStrategyType.FIRST.toStrategy()).isEqualTo(new FirstEventActivityStrategy());
}
@Test
public void testCreateLastEventStrategy() {
assertThat(ActivityStrategyType.LAST.toStrategy()).isEqualTo(LastEventActivityStrategy.getInstance());
}
@Test
public void testCreateFirstAndLastEventStrategy() {
assertThat(ActivityStrategyType.FIRST_AND_LAST.toStrategy()).isEqualTo(new FirstAndLastEventActivityStrategy());
}
}

42
common/transport/transport-api/src/test/java/org/thingsboard/server/common/transport/activity/strategy/AllEventsActivityStrategyTest.java

@ -0,0 +1,42 @@
/**
* Copyright © 2016-2024 The Thingsboard Authors
*
* Licensed under the Apache License, Version 2.0 (the "License");
* you may not use this file except in compliance with the License.
* You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package org.thingsboard.server.common.transport.activity.strategy;
import org.junit.jupiter.api.BeforeEach;
import org.junit.jupiter.api.Test;
import static org.junit.jupiter.api.Assertions.assertTrue;
public class AllEventsActivityStrategyTest {
private AllEventsActivityStrategy strategy;
@BeforeEach
public void setUp() {
strategy = AllEventsActivityStrategy.getInstance();
}
@Test
public void testOnActivity() {
assertTrue(strategy.onActivity(), "onActivity() should always return true.");
}
@Test
public void testOnReportingPeriodEnd() {
assertTrue(strategy.onReportingPeriodEnd(), "onReportingPeriodEnd() should always return true.");
}
}

53
common/transport/transport-api/src/test/java/org/thingsboard/server/common/transport/activity/strategy/FirstAndLastEventActivityStrategyTest.java

@ -0,0 +1,53 @@
/**
* Copyright © 2016-2024 The Thingsboard Authors
*
* Licensed under the Apache License, Version 2.0 (the "License");
* you may not use this file except in compliance with the License.
* You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package org.thingsboard.server.common.transport.activity.strategy;
import org.junit.jupiter.api.BeforeEach;
import org.junit.jupiter.api.Test;
import static org.junit.jupiter.api.Assertions.assertFalse;
import static org.junit.jupiter.api.Assertions.assertTrue;
public class FirstAndLastEventActivityStrategyTest {
private FirstAndLastEventActivityStrategy strategy;
@BeforeEach
public void setUp() {
strategy = new FirstAndLastEventActivityStrategy();
}
@Test
public void testOnActivity_FirstCall() {
assertTrue(strategy.onActivity(), "First call of onActivity() should return true.");
}
@Test
public void testOnActivity_SubsequentCalls() {
assertTrue(strategy.onActivity(), "First call of onActivity() should return true.");
assertFalse(strategy.onActivity(), "Subsequent calls of onActivity() should return false.");
}
@Test
public void testOnReportingPeriodEnd() {
assertTrue(strategy.onActivity(), "First call of onActivity() should return true.");
assertTrue(strategy.onReportingPeriodEnd(), "onReportingPeriodEnd() should always return true.");
assertTrue(strategy.onActivity(), "onActivity() should return true after onReportingPeriodEnd() for the next reporting period");
assertTrue(strategy.onReportingPeriodEnd(), "onReportingPeriodEnd() should always return true.");
}
}

52
common/transport/transport-api/src/test/java/org/thingsboard/server/common/transport/activity/strategy/FirstEventActivityStrategyTest.java

@ -0,0 +1,52 @@
/**
* Copyright © 2016-2024 The Thingsboard Authors
*
* Licensed under the Apache License, Version 2.0 (the "License");
* you may not use this file except in compliance with the License.
* You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package org.thingsboard.server.common.transport.activity.strategy;
import org.junit.jupiter.api.BeforeEach;
import org.junit.jupiter.api.Test;
import static org.junit.jupiter.api.Assertions.assertFalse;
import static org.junit.jupiter.api.Assertions.assertTrue;
public class FirstEventActivityStrategyTest {
private FirstEventActivityStrategy strategy;
@BeforeEach
public void setUp() {
strategy = new FirstEventActivityStrategy();
}
@Test
public void testOnActivity_FirstCall() {
assertTrue(strategy.onActivity(), "First call of onActivity() should return true.");
}
@Test
public void testOnActivity_SubsequentCalls() {
assertTrue(strategy.onActivity(), "First call of onActivity() should return true.");
assertFalse(strategy.onActivity(), "Subsequent calls of onActivity() should return false.");
}
@Test
public void testOnReportingPeriodEnd() {
assertTrue(strategy.onActivity(), "First call of onActivity() should return true.");
assertFalse(strategy.onReportingPeriodEnd(), "onReportingPeriodEnd() should always return false.");
assertTrue(strategy.onActivity(), "onActivity() should return true after onReportingPeriodEnd().");
assertFalse(strategy.onReportingPeriodEnd(), "onReportingPeriodEnd() should always return false.");
}
}

43
common/transport/transport-api/src/test/java/org/thingsboard/server/common/transport/activity/strategy/LastEventActivityStrategyTest.java

@ -0,0 +1,43 @@
/**
* Copyright © 2016-2024 The Thingsboard Authors
*
* Licensed under the Apache License, Version 2.0 (the "License");
* you may not use this file except in compliance with the License.
* You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package org.thingsboard.server.common.transport.activity.strategy;
import org.junit.jupiter.api.BeforeEach;
import org.junit.jupiter.api.Test;
import static org.junit.jupiter.api.Assertions.assertFalse;
import static org.junit.jupiter.api.Assertions.assertTrue;
public class LastEventActivityStrategyTest {
private LastEventActivityStrategy strategy;
@BeforeEach
public void setUp() {
strategy = LastEventActivityStrategy.getInstance();
}
@Test
public void testOnActivity() {
assertFalse(strategy.onActivity(), "onActivity() should always return false.");
}
@Test
public void testOnReportingPeriodEnd() {
assertTrue(strategy.onReportingPeriodEnd(), "onReportingPeriodEnd() should always return true.");
}
}

508
common/transport/transport-api/src/test/java/org/thingsboard/server/common/transport/service/TransportActivityManagerTest.java

@ -0,0 +1,508 @@
/**
* Copyright © 2016-2024 The Thingsboard Authors
*
* Licensed under the Apache License, Version 2.0 (the "License");
* you may not use this file except in compliance with the License.
* You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package org.thingsboard.server.common.transport.service;
import org.junit.jupiter.api.BeforeEach;
import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.extension.ExtendWith;
import org.junit.jupiter.params.ParameterizedTest;
import org.junit.jupiter.params.provider.Arguments;
import org.junit.jupiter.params.provider.EnumSource;
import org.junit.jupiter.params.provider.MethodSource;
import org.mockito.ArgumentCaptor;
import org.mockito.Mock;
import org.mockito.junit.jupiter.MockitoExtension;
import org.springframework.test.util.ReflectionTestUtils;
import org.thingsboard.server.common.transport.SessionMsgListener;
import org.thingsboard.server.common.transport.TransportServiceCallback;
import org.thingsboard.server.common.transport.activity.ActivityReportCallback;
import org.thingsboard.server.common.transport.activity.ActivityState;
import org.thingsboard.server.common.transport.activity.strategy.ActivityStrategy;
import org.thingsboard.server.common.transport.activity.strategy.ActivityStrategyType;
import org.thingsboard.server.gen.transport.TransportProtos;
import java.util.UUID;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.ConcurrentMap;
import java.util.stream.Stream;
import static org.assertj.core.api.Assertions.assertThat;
import static org.mockito.Mockito.doCallRealMethod;
import static org.mockito.Mockito.mock;
import static org.mockito.Mockito.never;
import static org.mockito.Mockito.verify;
import static org.mockito.Mockito.when;
import static org.thingsboard.server.common.transport.service.DefaultTransportService.SESSION_EVENT_MSG_CLOSED;
import static org.thingsboard.server.common.transport.service.DefaultTransportService.SESSION_EXPIRED_NOTIFICATION_PROTO;
@ExtendWith(MockitoExtension.class)
public class TransportActivityManagerTest {
private final UUID SESSION_ID = UUID.fromString("1306648a-9b26-11ee-b9d1-0242ac120002");
@Mock
private DefaultTransportService transportServiceMock;
private ConcurrentMap<UUID, SessionMetaData> sessions;
@BeforeEach
public void setup() {
sessions = new ConcurrentHashMap<>();
ReflectionTestUtils.setField(transportServiceMock, "sessions", sessions);
}
@Test
void givenKeyAndTimeToReportAndSessionExists_whenReportingActivity_thenShouldReportActivityWithSubscriptionsAndSessionInfoFromSession() {
// GIVEN
long expectedTime = 123L;
boolean expectedAttributesSubscription = true;
boolean expectedRPCSubscription = true;
TransportProtos.SessionInfoProto expectedSessionInfo = TransportProtos.SessionInfoProto.getDefaultInstance();
SessionMsgListener listenerMock = mock(SessionMsgListener.class);
SessionMetaData session = new SessionMetaData(expectedSessionInfo, TransportProtos.SessionType.ASYNC, listenerMock);
session.setSubscribedToAttributes(expectedAttributesSubscription);
session.setSubscribedToRPC(expectedRPCSubscription);
sessions.put(SESSION_ID, session);
ActivityReportCallback<UUID> callbackMock = mock(ActivityReportCallback.class);
TransportProtos.SessionInfoProto sessionInfo = TransportProtos.SessionInfoProto.newBuilder()
.setSessionIdMSB(SESSION_ID.getMostSignificantBits())
.setSessionIdLSB(SESSION_ID.getLeastSignificantBits())
.build();
doCallRealMethod().when(transportServiceMock).reportActivity(SESSION_ID, sessionInfo, expectedTime, callbackMock);
// WHEN
transportServiceMock.reportActivity(SESSION_ID, sessionInfo, expectedTime, callbackMock);
// THEN
ArgumentCaptor<TransportProtos.SessionInfoProto> sessionInfoCaptor = ArgumentCaptor.forClass(TransportProtos.SessionInfoProto.class);
ArgumentCaptor<TransportProtos.SubscriptionInfoProto> subscriptionInfoCaptor = ArgumentCaptor.forClass(TransportProtos.SubscriptionInfoProto.class);
ArgumentCaptor<TransportServiceCallback<Void>> callbackCaptor = ArgumentCaptor.forClass(TransportServiceCallback.class);
verify(transportServiceMock).process(sessionInfoCaptor.capture(), subscriptionInfoCaptor.capture(), callbackCaptor.capture());
assertThat(sessionInfoCaptor.getValue()).isEqualTo(expectedSessionInfo);
TransportProtos.SubscriptionInfoProto expectedSubscriptionInfo = TransportProtos.SubscriptionInfoProto.newBuilder()
.setAttributeSubscription(expectedAttributesSubscription)
.setRpcSubscription(expectedRPCSubscription)
.setLastActivityTime(expectedTime)
.build();
assertThat(subscriptionInfoCaptor.getValue()).isEqualTo(expectedSubscriptionInfo);
TransportServiceCallback<Void> queueCallback = callbackCaptor.getValue();
queueCallback.onSuccess(null);
verify(callbackMock).onSuccess(SESSION_ID, expectedTime);
var throwable = new Throwable();
queueCallback.onError(throwable);
verify(callbackMock).onFailure(SESSION_ID, throwable);
}
@Test
void givenKeyAndTimeToReportAndSessionDoesNotExist_whenReportingActivity_thenShouldReportActivityWithNoSubscriptionsAndPreviousSessionInfo() {
// GIVEN
long expectedTime = 123L;
boolean expectedAttributesSubscription = false;
boolean expectedRPCSubscription = false;
TransportProtos.SessionInfoProto expectedSessionInfo = TransportProtos.SessionInfoProto.newBuilder()
.setSessionIdMSB(SESSION_ID.getMostSignificantBits())
.setSessionIdLSB(SESSION_ID.getLeastSignificantBits())
.build();
ActivityReportCallback<UUID> callbackMock = mock(ActivityReportCallback.class);
doCallRealMethod().when(transportServiceMock).reportActivity(SESSION_ID, expectedSessionInfo, expectedTime, callbackMock);
// WHEN
transportServiceMock.reportActivity(SESSION_ID, expectedSessionInfo, expectedTime, callbackMock);
// THEN
ArgumentCaptor<TransportProtos.SessionInfoProto> sessionInfoCaptor = ArgumentCaptor.forClass(TransportProtos.SessionInfoProto.class);
ArgumentCaptor<TransportProtos.SubscriptionInfoProto> subscriptionInfoCaptor = ArgumentCaptor.forClass(TransportProtos.SubscriptionInfoProto.class);
ArgumentCaptor<TransportServiceCallback<Void>> callbackCaptor = ArgumentCaptor.forClass(TransportServiceCallback.class);
verify(transportServiceMock).process(sessionInfoCaptor.capture(), subscriptionInfoCaptor.capture(), callbackCaptor.capture());
assertThat(sessionInfoCaptor.getValue()).isEqualTo(expectedSessionInfo);
TransportProtos.SubscriptionInfoProto expectedSubscriptionInfo = TransportProtos.SubscriptionInfoProto.newBuilder()
.setAttributeSubscription(expectedAttributesSubscription)
.setRpcSubscription(expectedRPCSubscription)
.setLastActivityTime(expectedTime)
.build();
assertThat(subscriptionInfoCaptor.getValue()).isEqualTo(expectedSubscriptionInfo);
TransportServiceCallback<Void> queueCallback = callbackCaptor.getValue();
queueCallback.onSuccess(null);
verify(callbackMock).onSuccess(SESSION_ID, expectedTime);
var throwable = new Throwable();
queueCallback.onError(throwable);
verify(callbackMock).onFailure(SESSION_ID, throwable);
}
@Test
void givenActivityHappened_whenRecordActivity_thenShouldDelegateToOnActivity() {
// GIVEN
TransportProtos.SessionInfoProto sessionInfo = TransportProtos.SessionInfoProto.newBuilder()
.setSessionIdMSB(SESSION_ID.getMostSignificantBits())
.setSessionIdLSB(SESSION_ID.getLeastSignificantBits())
.build();
doCallRealMethod().when(transportServiceMock).recordActivity(sessionInfo);
when(transportServiceMock.toSessionId(sessionInfo)).thenReturn(SESSION_ID);
long expectedTime = 123L;
when(transportServiceMock.getCurrentTimeMillis()).thenReturn(expectedTime);
// WHEN
transportServiceMock.recordActivity(sessionInfo);
// THEN
verify(transportServiceMock).onActivity(SESSION_ID, expectedTime);
}
@Test
void givenKey_whenCreatingNewState_thenShouldCorrectlyCreateNewEmptyState() {
// GIVEN
TransportProtos.SessionInfoProto sessionInfo = TransportProtos.SessionInfoProto.newBuilder()
.setSessionIdMSB(SESSION_ID.getMostSignificantBits())
.setSessionIdLSB(SESSION_ID.getLeastSignificantBits())
.build();
sessions.put(SESSION_ID, new SessionMetaData(sessionInfo, TransportProtos.SessionType.ASYNC, null));
when(transportServiceMock.createNewState(SESSION_ID)).thenCallRealMethod();
ActivityState<TransportProtos.SessionInfoProto> expectedState = new ActivityState<>();
expectedState.setMetadata(sessionInfo);
// WHEN
ActivityState<TransportProtos.SessionInfoProto> actualState = transportServiceMock.createNewState(SESSION_ID);
// THEN
assertThat(actualState).isEqualTo(expectedState);
}
@ParameterizedTest
@EnumSource(ActivityStrategyType.class)
void givenDifferentReportingStrategies_whenGettingStrategy_thenShouldReturnCorrectStrategy(ActivityStrategyType reportingStrategyType) {
// GIVEN
doCallRealMethod().when(transportServiceMock).getStrategy();
ReflectionTestUtils.setField(transportServiceMock, "reportingStrategyType", reportingStrategyType);
// WHEN
ActivityStrategy actualStrategy = transportServiceMock.getStrategy();
// THEN
assertThat(actualStrategy).isEqualTo(reportingStrategyType.toStrategy());
}
@Test
void givenSessionDoesNotExist_whenUpdatingActivityState_thenShouldReturnNull() {
// GIVEN
TransportProtos.SessionInfoProto sessionInfo = TransportProtos.SessionInfoProto.newBuilder()
.setSessionIdMSB(SESSION_ID.getMostSignificantBits())
.setSessionIdLSB(SESSION_ID.getLeastSignificantBits())
.build();
ActivityState<TransportProtos.SessionInfoProto> state = new ActivityState<>();
state.setLastRecordedTime(123L);
state.setMetadata(sessionInfo);
when(transportServiceMock.updateState(SESSION_ID, state)).thenCallRealMethod();
// WHEN
ActivityState<TransportProtos.SessionInfoProto> updatedState = transportServiceMock.updateState(SESSION_ID, state);
// THEN
assertThat(updatedState).isNull();
}
@Test
void givenNoGwSessionId_whenUpdatingActivityState_thenShouldReturnSameInstanceWithUpdatedSessionInfo() {
// GIVEN
TransportProtos.SessionInfoProto sessionInfo = TransportProtos.SessionInfoProto.newBuilder()
.setSessionIdMSB(SESSION_ID.getMostSignificantBits())
.setSessionIdLSB(SESSION_ID.getLeastSignificantBits())
.build();
SessionMsgListener listenerMock = mock(SessionMsgListener.class);
sessions.put(SESSION_ID, new SessionMetaData(sessionInfo, TransportProtos.SessionType.ASYNC, listenerMock));
long lastRecordedTime = 123L;
ActivityState<TransportProtos.SessionInfoProto> state = new ActivityState<>();
state.setLastRecordedTime(lastRecordedTime);
state.setMetadata(TransportProtos.SessionInfoProto.getDefaultInstance());
when(transportServiceMock.updateState(SESSION_ID, state)).thenCallRealMethod();
// WHEN
ActivityState<TransportProtos.SessionInfoProto> updatedState = transportServiceMock.updateState(SESSION_ID, state);
// THEN
assertThat(updatedState).isSameAs(state);
assertThat(updatedState.getLastRecordedTime()).isEqualTo(lastRecordedTime);
assertThat(updatedState.getMetadata()).isEqualTo(sessionInfo);
}
@Test
void givenHasGwSessionIdButGwSessionIsNotNull_whenUpdatingActivityState_thenShouldReturnSameInstanceWithUpdatedSessionInfo() {
// GIVEN
var gwSessionId = UUID.fromString("19864038-9b48-11ee-b9d1-0242ac120002");
TransportProtos.SessionInfoProto sessionInfo = TransportProtos.SessionInfoProto.newBuilder()
.setSessionIdMSB(SESSION_ID.getMostSignificantBits())
.setSessionIdLSB(SESSION_ID.getLeastSignificantBits())
.setGwSessionIdMSB(gwSessionId.getMostSignificantBits())
.setGwSessionIdLSB(gwSessionId.getLeastSignificantBits())
.build();
SessionMsgListener listenerMock = mock(SessionMsgListener.class);
sessions.put(SESSION_ID, new SessionMetaData(sessionInfo, TransportProtos.SessionType.ASYNC, listenerMock));
long lastRecordedTime = 123L;
ActivityState<TransportProtos.SessionInfoProto> state = new ActivityState<>();
state.setLastRecordedTime(lastRecordedTime);
state.setMetadata(TransportProtos.SessionInfoProto.getDefaultInstance());
when(transportServiceMock.updateState(SESSION_ID, state)).thenCallRealMethod();
// WHEN
ActivityState<TransportProtos.SessionInfoProto> updatedState = transportServiceMock.updateState(SESSION_ID, state);
// THEN
assertThat(updatedState).isSameAs(state);
assertThat(updatedState.getLastRecordedTime()).isEqualTo(lastRecordedTime);
assertThat(updatedState.getMetadata()).isEqualTo(sessionInfo);
verify(transportServiceMock, never()).getLastRecordedTime(gwSessionId);
}
@Test
void givenHasGwSessionWithoutOverwriteEnabled_whenUpdatingActivityState_thenShouldReturnSameInstanceWithUpdatedSessionInfo() {
// GIVEN
var gwSessionId = UUID.fromString("19864038-9b48-11ee-b9d1-0242ac120002");
TransportProtos.SessionInfoProto gwSessionInfo = TransportProtos.SessionInfoProto.newBuilder()
.setSessionIdMSB(gwSessionId.getMostSignificantBits())
.setSessionIdLSB(gwSessionId.getLeastSignificantBits())
.build();
SessionMsgListener gwListenerMock = mock(SessionMsgListener.class);
sessions.put(gwSessionId, new SessionMetaData(gwSessionInfo, TransportProtos.SessionType.ASYNC, gwListenerMock));
TransportProtos.SessionInfoProto sessionInfo = TransportProtos.SessionInfoProto.newBuilder()
.setSessionIdMSB(SESSION_ID.getMostSignificantBits())
.setSessionIdLSB(SESSION_ID.getLeastSignificantBits())
.setGwSessionIdMSB(gwSessionId.getMostSignificantBits())
.setGwSessionIdLSB(gwSessionId.getLeastSignificantBits())
.build();
SessionMsgListener listenerMock = mock(SessionMsgListener.class);
sessions.put(SESSION_ID, new SessionMetaData(sessionInfo, TransportProtos.SessionType.ASYNC, listenerMock));
long lastRecordedTime = 123L;
ActivityState<TransportProtos.SessionInfoProto> state = new ActivityState<>();
state.setLastRecordedTime(lastRecordedTime);
state.setMetadata(TransportProtos.SessionInfoProto.getDefaultInstance());
when(transportServiceMock.updateState(SESSION_ID, state)).thenCallRealMethod();
// WHEN
ActivityState<TransportProtos.SessionInfoProto> updatedState = transportServiceMock.updateState(SESSION_ID, state);
// THEN
assertThat(updatedState).isSameAs(state);
assertThat(updatedState.getLastRecordedTime()).isEqualTo(lastRecordedTime);
assertThat(updatedState.getMetadata()).isEqualTo(sessionInfo);
verify(transportServiceMock, never()).getLastRecordedTime(gwSessionId);
}
@Test
void givenHasGwSessionWithOverwriteEnabledAndGwLastRecordedTimeIsGreater_whenUpdatingActivityState_thenShouldReturnSameInstanceWithUpdatedSessionInfoAndLastRecordedTime() {
// GIVEN
var gwSessionId = UUID.fromString("19864038-9b48-11ee-b9d1-0242ac120002");
TransportProtos.SessionInfoProto gwSessionInfo = TransportProtos.SessionInfoProto.newBuilder()
.setSessionIdMSB(gwSessionId.getMostSignificantBits())
.setSessionIdLSB(gwSessionId.getLeastSignificantBits())
.build();
SessionMsgListener gwListenerMock = mock(SessionMsgListener.class);
SessionMetaData gwSession = new SessionMetaData(gwSessionInfo, TransportProtos.SessionType.ASYNC, gwListenerMock);
gwSession.setOverwriteActivityTime(true);
sessions.put(gwSessionId, gwSession);
long gwLastRecordedTime = 500L;
when(transportServiceMock.getLastRecordedTime(gwSessionId)).thenReturn(gwLastRecordedTime);
TransportProtos.SessionInfoProto sessionInfo = TransportProtos.SessionInfoProto.newBuilder()
.setSessionIdMSB(SESSION_ID.getMostSignificantBits())
.setSessionIdLSB(SESSION_ID.getLeastSignificantBits())
.setGwSessionIdMSB(gwSessionId.getMostSignificantBits())
.setGwSessionIdLSB(gwSessionId.getLeastSignificantBits())
.build();
SessionMsgListener listenerMock = mock(SessionMsgListener.class);
sessions.put(SESSION_ID, new SessionMetaData(sessionInfo, TransportProtos.SessionType.ASYNC, listenerMock));
long lastRecordedTime = 123L;
ActivityState<TransportProtos.SessionInfoProto> state = new ActivityState<>();
state.setLastRecordedTime(lastRecordedTime);
state.setMetadata(TransportProtos.SessionInfoProto.getDefaultInstance());
when(transportServiceMock.updateState(SESSION_ID, state)).thenCallRealMethod();
// WHEN
ActivityState<TransportProtos.SessionInfoProto> updatedState = transportServiceMock.updateState(SESSION_ID, state);
// THEN
assertThat(updatedState).isSameAs(state);
assertThat(updatedState.getLastRecordedTime()).isEqualTo(gwLastRecordedTime);
assertThat(updatedState.getMetadata()).isEqualTo(sessionInfo);
}
@Test
void givenHasGwSessionWithOverwriteEnabledAndGwLastRecordedTimeIsLess_whenUpdatingActivityState_thenShouldReturnSameInstanceWithUpdatedSessionInfoOnly() {
// GIVEN
var gwSessionId = UUID.fromString("19864038-9b48-11ee-b9d1-0242ac120002");
TransportProtos.SessionInfoProto gwSessionInfo = TransportProtos.SessionInfoProto.newBuilder()
.setSessionIdMSB(gwSessionId.getMostSignificantBits())
.setSessionIdLSB(gwSessionId.getLeastSignificantBits())
.build();
SessionMsgListener gwListenerMock = mock(SessionMsgListener.class);
SessionMetaData gwSession = new SessionMetaData(gwSessionInfo, TransportProtos.SessionType.ASYNC, gwListenerMock);
gwSession.setOverwriteActivityTime(true);
sessions.put(gwSessionId, gwSession);
long gwLastRecordedTime = 100L;
when(transportServiceMock.getLastRecordedTime(gwSessionId)).thenReturn(gwLastRecordedTime);
TransportProtos.SessionInfoProto sessionInfo = TransportProtos.SessionInfoProto.newBuilder()
.setSessionIdMSB(SESSION_ID.getMostSignificantBits())
.setSessionIdLSB(SESSION_ID.getLeastSignificantBits())
.setGwSessionIdMSB(gwSessionId.getMostSignificantBits())
.setGwSessionIdLSB(gwSessionId.getLeastSignificantBits())
.build();
SessionMsgListener listenerMock = mock(SessionMsgListener.class);
sessions.put(SESSION_ID, new SessionMetaData(sessionInfo, TransportProtos.SessionType.ASYNC, listenerMock));
long lastRecordedTime = 123L;
ActivityState<TransportProtos.SessionInfoProto> state = new ActivityState<>();
state.setLastRecordedTime(lastRecordedTime);
state.setMetadata(TransportProtos.SessionInfoProto.getDefaultInstance());
when(transportServiceMock.updateState(SESSION_ID, state)).thenCallRealMethod();
// WHEN
ActivityState<TransportProtos.SessionInfoProto> updatedState = transportServiceMock.updateState(SESSION_ID, state);
// THEN
assertThat(updatedState).isSameAs(state);
assertThat(updatedState.getLastRecordedTime()).isEqualTo(lastRecordedTime);
assertThat(updatedState.getMetadata()).isEqualTo(sessionInfo);
}
@ParameterizedTest
@MethodSource("provideTestParamsForHasExpiredTrue")
public void givenExpiredLastRecordedTime_whenCheckingForExpiry_thenShouldReturnTrue(long currentTimeMillis, long lastRecordedTime, long sessionInactivityTimeout) {
// GIVEN
ReflectionTestUtils.setField(transportServiceMock, "sessionInactivityTimeout", sessionInactivityTimeout);
when(transportServiceMock.getCurrentTimeMillis()).thenReturn(currentTimeMillis);
when(transportServiceMock.hasExpired(lastRecordedTime)).thenCallRealMethod();
// WHEN
boolean hasExpired = transportServiceMock.hasExpired(lastRecordedTime);
// THEN
assertThat(hasExpired).isTrue();
}
private static Stream<Arguments> provideTestParamsForHasExpiredTrue() {
return Stream.of(
Arguments.of(10L, 0L, 9L),
Arguments.of(10L, 7L, 2L),
Arguments.of(10L, 8L, 1L),
Arguments.of(10000L, 5000L, 3000L)
);
}
@ParameterizedTest
@MethodSource("provideTestParamsForHasExpiredFalse")
public void givenNotExpiredLastRecordedTime_whenCheckingForExpiry_thenShouldReturnFalse(long currentTimeMillis, long lastRecordedTime, long sessionInactivityTimeout) {
// GIVEN
ReflectionTestUtils.setField(transportServiceMock, "sessionInactivityTimeout", sessionInactivityTimeout);
when(transportServiceMock.getCurrentTimeMillis()).thenReturn(currentTimeMillis);
when(transportServiceMock.hasExpired(lastRecordedTime)).thenCallRealMethod();
// WHEN
boolean hasExpired = transportServiceMock.hasExpired(lastRecordedTime);
// THEN
assertThat(hasExpired).isFalse();
}
private static Stream<Arguments> provideTestParamsForHasExpiredFalse() {
return Stream.of(
Arguments.of(10L, 9L, 2L),
Arguments.of(10L, 0L, 11L),
Arguments.of(10L, 8L, 3L),
Arguments.of(10000L, 8000L, 3000L)
);
}
@Test
void givenSessionExists_whenOnStateExpiryCalled_thenShouldPerformExpirationActions() {
// GIVEN
TransportProtos.SessionInfoProto sessionInfo = TransportProtos.SessionInfoProto.newBuilder()
.setSessionIdMSB(SESSION_ID.getMostSignificantBits())
.setSessionIdLSB(SESSION_ID.getLeastSignificantBits())
.build();
SessionMsgListener listenerMock = mock(SessionMsgListener.class);
sessions.put(SESSION_ID, new SessionMetaData(sessionInfo, TransportProtos.SessionType.ASYNC, listenerMock));
doCallRealMethod().when(transportServiceMock).onStateExpiry(SESSION_ID, sessionInfo);
// WHEN
transportServiceMock.onStateExpiry(SESSION_ID, sessionInfo);
// THEN
assertThat(sessions.containsKey(SESSION_ID)).isFalse();
verify(transportServiceMock).deregisterSession(sessionInfo);
verify(transportServiceMock).process(sessionInfo, SESSION_EVENT_MSG_CLOSED, null);
verify(listenerMock).onRemoteSessionCloseCommand(SESSION_ID, SESSION_EXPIRED_NOTIFICATION_PROTO);
}
@Test
void givenSessionDoesNotExist_whenOnStateExpiryCalled_thenShouldNotPerformExpirationActions() {
// GIVEN
TransportProtos.SessionInfoProto sessionInfo = TransportProtos.SessionInfoProto.newBuilder()
.setSessionIdMSB(SESSION_ID.getMostSignificantBits())
.setSessionIdLSB(SESSION_ID.getLeastSignificantBits())
.build();
doCallRealMethod().when(transportServiceMock).onStateExpiry(SESSION_ID, sessionInfo);
// WHEN
transportServiceMock.onStateExpiry(SESSION_ID, sessionInfo);
// THEN
assertThat(sessions.containsKey(SESSION_ID)).isFalse();
verify(transportServiceMock, never()).deregisterSession(sessionInfo);
verify(transportServiceMock, never()).process(sessionInfo, SESSION_EVENT_MSG_CLOSED, null);
}
}
Loading…
Cancel
Save