Browse Source

Additional improvements during testing and add tests for transport activity manager

pull/9980/head
Dmytro Skarzhynets 3 years ago
parent
commit
f42ea79726
  1. 58
      application/src/test/java/org/thingsboard/server/service/integration/IntegrationActivityManagerTest.java
  2. 2
      common/transport/coap/src/main/java/org/thingsboard/server/transport/coap/client/DefaultCoapClientContext.java
  3. 10
      common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportHandler.java
  4. 2
      common/transport/snmp/src/main/java/org/thingsboard/server/transport/snmp/service/SnmpTransportService.java
  5. 2
      common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/TransportService.java
  6. 3
      common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/activity/AbstractActivityManager.java
  7. 2
      common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/activity/ActivityManager.java
  8. 3
      common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/activity/strategy/AllEventsActivityStrategy.java
  9. 3
      common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/activity/strategy/FirstAndLastEventActivityStrategy.java
  10. 3
      common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/activity/strategy/FirstEventActivityStrategy.java
  11. 3
      common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/activity/strategy/LastEventActivityStrategy.java
  12. 40
      common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/service/DefaultTransportService.java
  13. 588
      common/transport/transport-api/src/test/java/org/thingsboard/server/common/transport/service/TransportActivityManagerTest.java

58
application/src/test/java/org/thingsboard/server/service/integration/IntegrationActivityManagerTest.java

@ -81,7 +81,7 @@ public class IntegrationActivityManagerTest {
private DefaultPlatformIntegrationService integrationServiceMock; private DefaultPlatformIntegrationService integrationServiceMock;
@Test @Test
void testReportActivity() { void givenKeyAndTimeToReport_whenReportingActivity_thenShouldCorrectlyReportActivity() {
// GIVEN // GIVEN
var key = new IntegrationActivityKey(TENANT_ID, DEVICE_ID); var key = new IntegrationActivityKey(TENANT_ID, DEVICE_ID);
@ -184,58 +184,54 @@ public class IntegrationActivityManagerTest {
@ParameterizedTest @ParameterizedTest
@MethodSource("provideTestParamsForCreateNewState") @MethodSource("provideTestParamsForCreateNewState")
void givenDifferentReportingStrategies_whenCreatingNewState_thenShouldCreateEmptyStateWithCorrectStrategy( void givenDifferentReportingStrategies_whenCreatingNewState_thenShouldCreateEmptyStateWithCorrectStrategy(
String reportingStrategyName, Class<ActivityStrategy> reportingStrategyClass String reportingStrategyName, ActivityStrategy reportingStrategy
) { ) {
// GIVEN // GIVEN
var key = new IntegrationActivityKey(TENANT_ID, DEVICE_ID); var key = new IntegrationActivityKey(TENANT_ID, DEVICE_ID);
when(integrationServiceMock.createNewState(key)).thenCallRealMethod(); when(integrationServiceMock.createNewState(key)).thenCallRealMethod();
ReflectionTestUtils.setField(integrationServiceMock, "reportingStrategyName", reportingStrategyName); ReflectionTestUtils.setField(integrationServiceMock, "reportingStrategyName", reportingStrategyName);
ActivityState<Void> expectedState = new ActivityState<>();
expectedState.setStrategy(reportingStrategy);
// WHEN // WHEN
ActivityState<Void> newState = integrationServiceMock.createNewState(key); ActivityState<Void> actualState = integrationServiceMock.createNewState(key);
// THEN // THEN
assertThat(newState).isNotNull(); assertThat(actualState).isEqualTo(expectedState);
assertThat(newState.getLastRecordedTime()).isEqualTo(0L);
assertThat(newState.getLastReportedTime()).isEqualTo(0L);
assertThat(newState.getStrategy()).isInstanceOf(reportingStrategyClass);
} }
private static Stream<Arguments> provideTestParamsForCreateNewState() { private static Stream<Arguments> provideTestParamsForCreateNewState() {
return Stream.of( return Stream.of(
Arguments.of("ALL", AllEventsActivityStrategy.class), Arguments.of("ALL", new AllEventsActivityStrategy()),
Arguments.of("FIRST", FirstEventActivityStrategy.class), Arguments.of("FIRST", new FirstEventActivityStrategy()),
Arguments.of("LAST", LastEventActivityStrategy.class), Arguments.of("LAST", new LastEventActivityStrategy()),
Arguments.of("FIRST_AND_LAST", FirstAndLastEventActivityStrategy.class) Arguments.of("FIRST_AND_LAST", new FirstAndLastEventActivityStrategy())
); );
} }
@Test @Test
void givenActivityState_whenUpdatingActivityState_thenShouldReturnSameInstanceWithNoChanges() { void givenActivityState_whenUpdatingActivityState_thenShouldReturnSameInstanceWithNoInteractions() {
// GIVEN // GIVEN
var key = new IntegrationActivityKey(TENANT_ID, DEVICE_ID); var key = new IntegrationActivityKey(TENANT_ID, DEVICE_ID);
long expectedLastRecordedTime = 123L; ActivityStrategy strategySpy = spy(new FirstEventActivityStrategy());
long expectedLastReportedTime = 312L;
ActivityStrategy expectedStrategy = spy(new FirstEventActivityStrategy());
ActivityState<Void> expectedState = new ActivityState<>(); ActivityState<Void> state = new ActivityState<>();
expectedState.setLastRecordedTime(expectedLastRecordedTime); state.setLastRecordedTime(123L);
expectedState.setLastReportedTime(expectedLastReportedTime); state.setLastReportedTime(312L);
expectedState.setStrategy(expectedStrategy); state.setStrategy(strategySpy);
ActivityState<Void> stateSpy = spy(state);
when(integrationServiceMock.updateState(key, expectedState)).thenCallRealMethod(); when(integrationServiceMock.updateState(key, state)).thenCallRealMethod();
// WHEN // WHEN
ActivityState<Void> actualNewState = integrationServiceMock.updateState(key, expectedState); ActivityState<Void> updatedState = integrationServiceMock.updateState(key, state);
// THEN // THEN
assertThat(actualNewState).isSameAs(expectedState); assertThat(updatedState).isSameAs(state);
assertThat(actualNewState.getLastRecordedTime()).isEqualTo(expectedLastRecordedTime); verifyNoInteractions(stateSpy);
assertThat(actualNewState.getLastReportedTime()).isEqualTo(expectedLastReportedTime); verifyNoInteractions(strategySpy);
assertThat(actualNewState.getStrategy()).isSameAs(expectedStrategy);
verifyNoInteractions(expectedStrategy);
assertThat(actualNewState.getMetadata()).isNull();
} }
@ParameterizedTest @ParameterizedTest
@ -258,7 +254,8 @@ public class IntegrationActivityManagerTest {
return Stream.of( return Stream.of(
Arguments.of(10L, 0L, 9L), Arguments.of(10L, 0L, 9L),
Arguments.of(10L, 7L, 2L), Arguments.of(10L, 7L, 2L),
Arguments.of(10L, 8L, 1L) Arguments.of(10L, 8L, 1L),
Arguments.of(10000L, 5000L, 3000L)
); );
} }
@ -282,12 +279,13 @@ public class IntegrationActivityManagerTest {
return Stream.of( return Stream.of(
Arguments.of(10L, 9L, 2L), Arguments.of(10L, 9L, 2L),
Arguments.of(10L, 0L, 11L), Arguments.of(10L, 0L, 11L),
Arguments.of(10L, 8L, 3L) Arguments.of(10L, 8L, 3L),
Arguments.of(10000L, 8000L, 3000L)
); );
} }
@Test @Test
void givenKeyAndVoidMetadata_whenOnStateExpiryCalled_thenShouldDoNothing() { void givenKeyAndMetadata_whenOnStateExpiryCalled_thenShouldDoNothing() {
// GIVEN // GIVEN
var key = new IntegrationActivityKey(TENANT_ID, DEVICE_ID); var key = new IntegrationActivityKey(TENANT_ID, DEVICE_ID);
doCallRealMethod().when(integrationServiceMock).onStateExpiry(key, null); doCallRealMethod().when(integrationServiceMock).onStateExpiry(key, null);

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() { public void reportActivity() {
for (TbCoapClientState state : clients.values()) { for (TbCoapClientState state : clients.values()) {
if (state.getSession() != null) { 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: case PINGREQ:
if (checkConnected(ctx, msg)) { if (checkConnected(ctx, msg)) {
ctx.writeAndFlush(new MqttMessage(new MqttFixedHeader(PINGRESP, false, AT_MOST_ONCE, false, 0))); ctx.writeAndFlush(new MqttMessage(new MqttFixedHeader(PINGRESP, false, AT_MOST_ONCE, false, 0)));
transportService.reportActivity(deviceSessionCtx.getSessionInfo()); transportService.recordActivity(deviceSessionCtx.getSessionInfo());
} }
break; break;
case DISCONNECT: case DISCONNECT:
@ -350,7 +350,7 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement
if (topicName.startsWith(MqttTopics.BASE_GATEWAY_API_TOPIC)) { if (topicName.startsWith(MqttTopics.BASE_GATEWAY_API_TOPIC)) {
if (gatewaySessionHandler != null) { if (gatewaySessionHandler != null) {
handleGatewayPublishMsg(ctx, topicName, msgId, mqttMsg); handleGatewayPublishMsg(ctx, topicName, msgId, mqttMsg);
transportService.reportActivity(deviceSessionCtx.getSessionInfo()); transportService.recordActivity(deviceSessionCtx.getSessionInfo());
} else { } else {
log.error("[gatewaySessionHandler] is null, [{}] Failed to process publish msg [{}][{}]", sessionId, topicName, msgId); 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)); transportService.process(deviceSessionCtx.getSessionInfo(), getAttributeMsg, getPubAckCallback(ctx, msgId, getAttributeMsg));
attrReqTopicType = TopicType.V2; attrReqTopicType = TopicType.V2;
} else { } else {
transportService.reportActivity(deviceSessionCtx.getSessionInfo()); transportService.recordActivity(deviceSessionCtx.getSessionInfo());
ack(ctx, msgId, ReturnCode.TOPIC_NAME_INVALID); ack(ctx, msgId, ReturnCode.TOPIC_NAME_INVALID);
} }
} catch (AdaptorException e) { } catch (AdaptorException e) {
@ -796,7 +796,7 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement
} }
} }
if (!activityReported) { if (!activityReported) {
transportService.reportActivity(deviceSessionCtx.getSessionInfo()); transportService.recordActivity(deviceSessionCtx.getSessionInfo());
} }
ctx.writeAndFlush(createSubAckMessage(mqttMsg.variableHeader().messageId(), grantedQoSList)); ctx.writeAndFlush(createSubAckMessage(mqttMsg.variableHeader().messageId(), grantedQoSList));
} }
@ -894,7 +894,7 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement
} }
} }
if (!activityReported) { if (!activityReported) {
transportService.reportActivity(deviceSessionCtx.getSessionInfo()); transportService.recordActivity(deviceSessionCtx.getSessionInfo());
} }
ctx.writeAndFlush(createUnSubAckMessage(mqttMsg.variableHeader().messageId(), unSubResults)); 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) { 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); 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); void lifecycleEvent(TenantId tenantId, DeviceId deviceId, ComponentLifecycleEvent eventType, boolean success, Throwable error);

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

@ -129,7 +129,8 @@ public abstract class AbstractActivityManager<Key, Metadata> implements Activity
} }
} }
protected long getLastRecordedTime(Key key) { @Override
public long getLastRecordedTime(Key key) {
ActivityState<Metadata> state = states.get(key); ActivityState<Metadata> state = states.get(key);
return state == null ? 0L : state.getLastRecordedTime(); return state == null ? 0L : state.getLastRecordedTime();
} }

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

@ -36,4 +36,6 @@ public interface ActivityManager<Key> {
void onActivity(Key key); void onActivity(Key key);
long getLastRecordedTime(Key key);
} }

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

@ -30,6 +30,9 @@
*/ */
package org.thingsboard.server.common.transport.activity.strategy; package org.thingsboard.server.common.transport.activity.strategy;
import lombok.EqualsAndHashCode;
@EqualsAndHashCode
public class AllEventsActivityStrategy implements ActivityStrategy { public class AllEventsActivityStrategy implements ActivityStrategy {
@Override @Override

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

@ -30,6 +30,9 @@
*/ */
package org.thingsboard.server.common.transport.activity.strategy; package org.thingsboard.server.common.transport.activity.strategy;
import lombok.EqualsAndHashCode;
@EqualsAndHashCode
public class FirstAndLastEventActivityStrategy implements ActivityStrategy { public class FirstAndLastEventActivityStrategy implements ActivityStrategy {
private boolean firstEventReceived; private boolean firstEventReceived;

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

@ -30,6 +30,9 @@
*/ */
package org.thingsboard.server.common.transport.activity.strategy; package org.thingsboard.server.common.transport.activity.strategy;
import lombok.EqualsAndHashCode;
@EqualsAndHashCode
public class FirstEventActivityStrategy implements ActivityStrategy { public class FirstEventActivityStrategy implements ActivityStrategy {
private boolean firstEventReceived; private boolean firstEventReceived;

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

@ -30,6 +30,9 @@
*/ */
package org.thingsboard.server.common.transport.activity.strategy; package org.thingsboard.server.common.transport.activity.strategy;
import lombok.EqualsAndHashCode;
@EqualsAndHashCode
public class LastEventActivityStrategy implements ActivityStrategy { public class LastEventActivityStrategy implements ActivityStrategy {
@Override @Override

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

@ -558,7 +558,7 @@ public class DefaultTransportService extends AbstractActivityManager<UUID, Trans
@Override @Override
public void process(TransportProtos.SessionInfoProto sessionInfo, TransportProtos.SessionEventMsg msg, TransportServiceCallback<Void> callback) { public void process(TransportProtos.SessionInfoProto sessionInfo, TransportProtos.SessionEventMsg msg, TransportServiceCallback<Void> callback) {
if (checkLimits(sessionInfo, msg, callback)) { if (checkLimits(sessionInfo, msg, callback)) {
reportActivityInternal(sessionInfo); recordActivityInternal(sessionInfo);
sendToDeviceActor(sessionInfo, TransportToDeviceActorMsg.newBuilder().setSessionInfo(sessionInfo) sendToDeviceActor(sessionInfo, TransportToDeviceActorMsg.newBuilder().setSessionInfo(sessionInfo)
.setSessionEvent(msg).build(), callback); .setSessionEvent(msg).build(), callback);
} }
@ -578,7 +578,7 @@ public class DefaultTransportService extends AbstractActivityManager<UUID, Trans
} }
} }
reportActivityInternal(sessionInfo); recordActivityInternal(sessionInfo);
sendToDeviceActor(sessionInfo, msg, callback); sendToDeviceActor(sessionInfo, msg, callback);
} }
} }
@ -595,7 +595,7 @@ public class DefaultTransportService extends AbstractActivityManager<UUID, Trans
dataPoints += tsKv.getKvCount(); dataPoints += tsKv.getKvCount();
} }
if (checkLimits(sessionInfo, msg, callback, dataPoints)) { if (checkLimits(sessionInfo, msg, callback, dataPoints)) {
reportActivityInternal(sessionInfo); recordActivityInternal(sessionInfo);
TenantId tenantId = getTenantId(sessionInfo); TenantId tenantId = getTenantId(sessionInfo);
DeviceId deviceId = new DeviceId(new UUID(sessionInfo.getDeviceIdMSB(), sessionInfo.getDeviceIdLSB())); DeviceId deviceId = new DeviceId(new UUID(sessionInfo.getDeviceIdMSB(), sessionInfo.getDeviceIdLSB()));
CustomerId customerId = getCustomerId(sessionInfo); CustomerId customerId = getCustomerId(sessionInfo);
@ -619,7 +619,7 @@ public class DefaultTransportService extends AbstractActivityManager<UUID, Trans
@Override @Override
public void process(TransportProtos.SessionInfoProto sessionInfo, TransportProtos.PostAttributeMsg msg, TbMsgMetaData md, TransportServiceCallback<Void> callback) { public void process(TransportProtos.SessionInfoProto sessionInfo, TransportProtos.PostAttributeMsg msg, TbMsgMetaData md, TransportServiceCallback<Void> callback) {
if (checkLimits(sessionInfo, msg, callback, msg.getKvCount())) { if (checkLimits(sessionInfo, msg, callback, msg.getKvCount())) {
reportActivityInternal(sessionInfo); recordActivityInternal(sessionInfo);
TenantId tenantId = getTenantId(sessionInfo); TenantId tenantId = getTenantId(sessionInfo);
DeviceId deviceId = new DeviceId(new UUID(sessionInfo.getDeviceIdMSB(), sessionInfo.getDeviceIdLSB())); DeviceId deviceId = new DeviceId(new UUID(sessionInfo.getDeviceIdMSB(), sessionInfo.getDeviceIdLSB()));
JsonObject json = JsonUtils.getJsonObject(msg.getKvList()); JsonObject json = JsonUtils.getJsonObject(msg.getKvList());
@ -639,7 +639,7 @@ public class DefaultTransportService extends AbstractActivityManager<UUID, Trans
@Override @Override
public void process(TransportProtos.SessionInfoProto sessionInfo, TransportProtos.GetAttributeRequestMsg msg, TransportServiceCallback<Void> callback) { public void process(TransportProtos.SessionInfoProto sessionInfo, TransportProtos.GetAttributeRequestMsg msg, TransportServiceCallback<Void> callback) {
if (checkLimits(sessionInfo, msg, callback)) { if (checkLimits(sessionInfo, msg, callback)) {
reportActivityInternal(sessionInfo); recordActivityInternal(sessionInfo);
sendToDeviceActor(sessionInfo, TransportToDeviceActorMsg.newBuilder().setSessionInfo(sessionInfo) sendToDeviceActor(sessionInfo, TransportToDeviceActorMsg.newBuilder().setSessionInfo(sessionInfo)
.setGetAttributes(msg).build(), new ApiStatsProxyCallback<>(getTenantId(sessionInfo), getCustomerId(sessionInfo), 1, callback)); .setGetAttributes(msg).build(), new ApiStatsProxyCallback<>(getTenantId(sessionInfo), getCustomerId(sessionInfo), 1, callback));
} }
@ -652,7 +652,7 @@ public class DefaultTransportService extends AbstractActivityManager<UUID, Trans
if (sessionMetaData != null) { if (sessionMetaData != null) {
sessionMetaData.setSubscribedToAttributes(!msg.getUnsubscribe()); sessionMetaData.setSubscribedToAttributes(!msg.getUnsubscribe());
} }
reportActivityInternal(sessionInfo); recordActivityInternal(sessionInfo);
sendToDeviceActor(sessionInfo, TransportToDeviceActorMsg.newBuilder().setSessionInfo(sessionInfo).setSubscribeToAttributes(msg).build(), sendToDeviceActor(sessionInfo, TransportToDeviceActorMsg.newBuilder().setSessionInfo(sessionInfo).setSubscribeToAttributes(msg).build(),
new ApiStatsProxyCallback<>(getTenantId(sessionInfo), getCustomerId(sessionInfo), 1, callback)); new ApiStatsProxyCallback<>(getTenantId(sessionInfo), getCustomerId(sessionInfo), 1, callback));
} }
@ -665,7 +665,7 @@ public class DefaultTransportService extends AbstractActivityManager<UUID, Trans
if (sessionMetaData != null) { if (sessionMetaData != null) {
sessionMetaData.setSubscribedToRPC(!msg.getUnsubscribe()); sessionMetaData.setSubscribedToRPC(!msg.getUnsubscribe());
} }
reportActivityInternal(sessionInfo); recordActivityInternal(sessionInfo);
sendToDeviceActor(sessionInfo, TransportToDeviceActorMsg.newBuilder().setSessionInfo(sessionInfo).setSubscribeToRPC(msg).build(), sendToDeviceActor(sessionInfo, TransportToDeviceActorMsg.newBuilder().setSessionInfo(sessionInfo).setSubscribeToRPC(msg).build(),
new ApiStatsProxyCallback<>(getTenantId(sessionInfo), getCustomerId(sessionInfo), 1, callback)); new ApiStatsProxyCallback<>(getTenantId(sessionInfo), getCustomerId(sessionInfo), 1, callback));
} }
@ -674,7 +674,7 @@ public class DefaultTransportService extends AbstractActivityManager<UUID, Trans
@Override @Override
public void process(TransportProtos.SessionInfoProto sessionInfo, TransportProtos.ToDeviceRpcResponseMsg msg, TransportServiceCallback<Void> callback) { public void process(TransportProtos.SessionInfoProto sessionInfo, TransportProtos.ToDeviceRpcResponseMsg msg, TransportServiceCallback<Void> callback) {
if (checkLimits(sessionInfo, msg, callback)) { if (checkLimits(sessionInfo, msg, callback)) {
reportActivityInternal(sessionInfo); recordActivityInternal(sessionInfo);
sendToDeviceActor(sessionInfo, TransportToDeviceActorMsg.newBuilder().setSessionInfo(sessionInfo).setToDeviceRPCCallResponse(msg).build(), sendToDeviceActor(sessionInfo, TransportToDeviceActorMsg.newBuilder().setSessionInfo(sessionInfo).setToDeviceRPCCallResponse(msg).build(),
new ApiStatsProxyCallback<>(getTenantId(sessionInfo), getCustomerId(sessionInfo), 1, callback)); new ApiStatsProxyCallback<>(getTenantId(sessionInfo), getCustomerId(sessionInfo), 1, callback));
} }
@ -683,7 +683,7 @@ public class DefaultTransportService extends AbstractActivityManager<UUID, Trans
@Override @Override
public void notifyAboutUplink(TransportProtos.SessionInfoProto sessionInfo, TransportProtos.UplinkNotificationMsg msg, TransportServiceCallback<Void> callback) { public void notifyAboutUplink(TransportProtos.SessionInfoProto sessionInfo, TransportProtos.UplinkNotificationMsg msg, TransportServiceCallback<Void> callback) {
if (checkLimits(sessionInfo, msg, callback)) { if (checkLimits(sessionInfo, msg, callback)) {
reportActivityInternal(sessionInfo); recordActivityInternal(sessionInfo);
sendToDeviceActor(sessionInfo, TransportToDeviceActorMsg.newBuilder().setSessionInfo(sessionInfo).setUplinkNotificationMsg(msg).build(), callback); sendToDeviceActor(sessionInfo, TransportToDeviceActorMsg.newBuilder().setSessionInfo(sessionInfo).setUplinkNotificationMsg(msg).build(), callback);
} }
} }
@ -704,7 +704,7 @@ public class DefaultTransportService extends AbstractActivityManager<UUID, Trans
if (checkLimits(sessionInfo, responseMsg, callback)) { if (checkLimits(sessionInfo, responseMsg, callback)) {
if (reportActivity) { if (reportActivity) {
reportActivityInternal(sessionInfo); recordActivityInternal(sessionInfo);
} }
sendToDeviceActor(sessionInfo, TransportToDeviceActorMsg.newBuilder().setSessionInfo(sessionInfo).setRpcResponseStatusMsg(responseMsg).build(), sendToDeviceActor(sessionInfo, TransportToDeviceActorMsg.newBuilder().setSessionInfo(sessionInfo).setRpcResponseStatusMsg(responseMsg).build(),
new ApiStatsProxyCallback<>(getTenantId(sessionInfo), getCustomerId(sessionInfo), 1, TransportServiceCallback.EMPTY)); new ApiStatsProxyCallback<>(getTenantId(sessionInfo), getCustomerId(sessionInfo), 1, TransportServiceCallback.EMPTY));
@ -736,7 +736,7 @@ public class DefaultTransportService extends AbstractActivityManager<UUID, Trans
@Override @Override
public void process(TransportProtos.SessionInfoProto sessionInfo, TransportProtos.ToServerRpcRequestMsg msg, TransportServiceCallback<Void> callback) { public void process(TransportProtos.SessionInfoProto sessionInfo, TransportProtos.ToServerRpcRequestMsg msg, TransportServiceCallback<Void> callback) {
if (checkLimits(sessionInfo, msg, callback)) { if (checkLimits(sessionInfo, msg, callback)) {
reportActivityInternal(sessionInfo); recordActivityInternal(sessionInfo);
UUID sessionId = toSessionId(sessionInfo); UUID sessionId = toSessionId(sessionInfo);
TenantId tenantId = getTenantId(sessionInfo); TenantId tenantId = getTenantId(sessionInfo);
DeviceId deviceId = getDeviceId(sessionInfo); DeviceId deviceId = getDeviceId(sessionInfo);
@ -761,7 +761,7 @@ public class DefaultTransportService extends AbstractActivityManager<UUID, Trans
@Override @Override
public void process(TransportProtos.SessionInfoProto sessionInfo, TransportProtos.ClaimDeviceMsg msg, TransportServiceCallback<Void> callback) { public void process(TransportProtos.SessionInfoProto sessionInfo, TransportProtos.ClaimDeviceMsg msg, TransportServiceCallback<Void> callback) {
if (checkLimits(sessionInfo, msg, callback)) { if (checkLimits(sessionInfo, msg, callback)) {
reportActivityInternal(sessionInfo); recordActivityInternal(sessionInfo);
sendToDeviceActor(sessionInfo, TransportToDeviceActorMsg.newBuilder().setSessionInfo(sessionInfo) sendToDeviceActor(sessionInfo, TransportToDeviceActorMsg.newBuilder().setSessionInfo(sessionInfo)
.setClaimDevice(msg).build(), callback); .setClaimDevice(msg).build(), callback);
} }
@ -780,11 +780,11 @@ public class DefaultTransportService extends AbstractActivityManager<UUID, Trans
} }
@Override @Override
public void reportActivity(TransportProtos.SessionInfoProto sessionInfo) { public void recordActivity(TransportProtos.SessionInfoProto sessionInfo) {
reportActivityInternal(sessionInfo); recordActivityInternal(sessionInfo);
} }
private void reportActivityInternal(TransportProtos.SessionInfoProto sessionInfo) { private void recordActivityInternal(TransportProtos.SessionInfoProto sessionInfo) {
onActivity(toSessionId(sessionInfo)); onActivity(toSessionId(sessionInfo));
} }
@ -810,7 +810,7 @@ public class DefaultTransportService extends AbstractActivityManager<UUID, Trans
state.setMetadata(session.getSessionInfo()); state.setMetadata(session.getSessionInfo());
var sessionInfo = state.getMetadata(); var sessionInfo = state.getMetadata();
if (sessionInfo.getGwSessionIdMSB() == 0L || sessionInfo.getCustomerIdLSB() == 0L) { if (sessionInfo.getGwSessionIdMSB() == 0L || sessionInfo.getGwSessionIdLSB() == 0L) {
return state; return state;
} }
@ -831,7 +831,11 @@ public class DefaultTransportService extends AbstractActivityManager<UUID, Trans
@Override @Override
protected boolean hasExpired(long lastRecordedTime) { protected boolean hasExpired(long lastRecordedTime) {
return (System.currentTimeMillis() - sessionInactivityTimeout) > lastRecordedTime; return (getCurrentTimeMillis() - sessionInactivityTimeout) > lastRecordedTime;
}
long getCurrentTimeMillis() {
return System.currentTimeMillis();
} }
@Override @Override
@ -937,7 +941,7 @@ public class DefaultTransportService extends AbstractActivityManager<UUID, Trans
} }
TransportProtos.PostTelemetryMsg.Builder request = TransportProtos.PostTelemetryMsg.newBuilder(); TransportProtos.PostTelemetryMsg.Builder request = TransportProtos.PostTelemetryMsg.newBuilder();
TransportProtos.TsKvListProto.Builder builder = TransportProtos.TsKvListProto.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() builder.addKv(TransportProtos.KeyValueProto.newBuilder()
.setKey("transportLog") .setKey("transportLog")
.setType(TransportProtos.KeyValueType.STRING_V) .setType(TransportProtos.KeyValueType.STRING_V)

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

@ -0,0 +1,588 @@
/**
* ThingsBoard, Inc. ("COMPANY") CONFIDENTIAL
*
* Copyright © 2016-2023 ThingsBoard, Inc. All Rights Reserved.
*
* NOTICE: All information contained herein is, and remains
* the property of ThingsBoard, Inc. and its suppliers,
* if any. The intellectual and technical concepts contained
* herein are proprietary to ThingsBoard, Inc.
* and its suppliers and may be covered by U.S. and Foreign Patents,
* patents in process, and are protected by trade secret or copyright law.
*
* Dissemination of this information or reproduction of this material is strictly forbidden
* unless prior written permission is obtained from COMPANY.
*
* Access to the source code contained herein is hereby forbidden to anyone except current COMPANY employees,
* managers or contractors who have executed Confidentiality and Non-disclosure agreements
* explicitly covering such access.
*
* The copyright notice above does not evidence any actual or intended publication
* or disclosure of this source code, which includes
* information that is confidential and/or proprietary, and is a trade secret, of COMPANY.
* ANY REPRODUCTION, MODIFICATION, DISTRIBUTION, PUBLIC PERFORMANCE,
* OR PUBLIC DISPLAY OF OR THROUGH USE OF THIS SOURCE CODE WITHOUT
* THE EXPRESS WRITTEN CONSENT OF COMPANY IS STRICTLY PROHIBITED,
* AND IN VIOLATION OF APPLICABLE LAWS AND INTERNATIONAL TREATIES.
* THE RECEIPT OR POSSESSION OF THIS SOURCE CODE AND/OR RELATED INFORMATION
* DOES NOT CONVEY OR IMPLY ANY RIGHTS TO REPRODUCE, DISCLOSE OR DISTRIBUTE ITS CONTENTS,
* OR TO MANUFACTURE, USE, OR SELL ANYTHING THAT IT MAY DESCRIBE, IN WHOLE OR IN PART.
*/
package org.thingsboard.server.common.transport.service;
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.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.AllEventsActivityStrategy;
import org.thingsboard.server.common.transport.activity.strategy.FirstAndLastEventActivityStrategy;
import org.thingsboard.server.common.transport.activity.strategy.FirstEventActivityStrategy;
import org.thingsboard.server.common.transport.activity.strategy.LastEventActivityStrategy;
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.spy;
import static org.mockito.Mockito.verify;
import static org.mockito.Mockito.verifyNoInteractions;
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);
}
// @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: [{}].", name, 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);
// }
// });
// }
@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);
// WHEN
transportServiceMock.recordActivity(sessionInfo);
// THEN
verify(transportServiceMock).onActivity(SESSION_ID);
}
@ParameterizedTest
@MethodSource("provideTestParamsForCreateNewState")
void givenDifferentReportingStrategies_whenCreatingNewState_thenShouldCreateEmptyStateWithCorrectStrategy(
String reportingStrategyName, ActivityStrategy reportingStrategy
) {
// 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));
ReflectionTestUtils.setField(transportServiceMock, "reportingStrategyName", reportingStrategyName);
when(transportServiceMock.createNewState(SESSION_ID)).thenCallRealMethod();
ActivityState<TransportProtos.SessionInfoProto> expectedState = new ActivityState<>();
expectedState.setStrategy(reportingStrategy);
expectedState.setMetadata(sessionInfo);
// WHEN
ActivityState<TransportProtos.SessionInfoProto> actualState = transportServiceMock.createNewState(SESSION_ID);
// THEN
assertThat(actualState).isEqualTo(expectedState);
}
private static Stream<Arguments> provideTestParamsForCreateNewState() {
return Stream.of(
Arguments.of("ALL", new AllEventsActivityStrategy()),
Arguments.of("FIRST", new FirstEventActivityStrategy()),
Arguments.of("LAST", new LastEventActivityStrategy()),
Arguments.of("FIRST_AND_LAST", new FirstAndLastEventActivityStrategy())
);
}
@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.setLastReportedTime(312L);
state.setMetadata(sessionInfo);
state.setStrategy(new FirstEventActivityStrategy());
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;
long lastReportedTime = 312L;
ActivityStrategy strategySpy = spy(new FirstEventActivityStrategy());
ActivityState<TransportProtos.SessionInfoProto> state = new ActivityState<>();
state.setLastRecordedTime(lastRecordedTime);
state.setLastReportedTime(lastReportedTime);
state.setMetadata(TransportProtos.SessionInfoProto.getDefaultInstance());
state.setStrategy(strategySpy);
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.getLastReportedTime()).isEqualTo(lastReportedTime);
assertThat(updatedState.getMetadata()).isEqualTo(sessionInfo);
assertThat(updatedState.getStrategy()).isEqualTo(strategySpy);
verifyNoInteractions(strategySpy);
}
@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;
long lastReportedTime = 312L;
ActivityStrategy strategySpy = spy(new FirstEventActivityStrategy());
ActivityState<TransportProtos.SessionInfoProto> state = new ActivityState<>();
state.setLastRecordedTime(lastRecordedTime);
state.setLastReportedTime(lastReportedTime);
state.setMetadata(TransportProtos.SessionInfoProto.getDefaultInstance());
state.setStrategy(strategySpy);
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.getLastReportedTime()).isEqualTo(lastReportedTime);
assertThat(updatedState.getMetadata()).isEqualTo(sessionInfo);
assertThat(updatedState.getStrategy()).isEqualTo(strategySpy);
verifyNoInteractions(strategySpy);
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;
long lastReportedTime = 312L;
ActivityStrategy strategySpy = spy(new FirstEventActivityStrategy());
ActivityState<TransportProtos.SessionInfoProto> state = new ActivityState<>();
state.setLastRecordedTime(lastRecordedTime);
state.setLastReportedTime(lastReportedTime);
state.setMetadata(TransportProtos.SessionInfoProto.getDefaultInstance());
state.setStrategy(strategySpy);
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.getLastReportedTime()).isEqualTo(lastReportedTime);
assertThat(updatedState.getMetadata()).isEqualTo(sessionInfo);
assertThat(updatedState.getStrategy()).isEqualTo(strategySpy);
verifyNoInteractions(strategySpy);
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;
long lastReportedTime = 312L;
ActivityStrategy strategySpy = spy(new FirstEventActivityStrategy());
ActivityState<TransportProtos.SessionInfoProto> state = new ActivityState<>();
state.setLastRecordedTime(lastRecordedTime);
state.setLastReportedTime(lastReportedTime);
state.setMetadata(TransportProtos.SessionInfoProto.getDefaultInstance());
state.setStrategy(strategySpy);
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.getLastReportedTime()).isEqualTo(lastReportedTime);
assertThat(updatedState.getMetadata()).isEqualTo(sessionInfo);
assertThat(updatedState.getStrategy()).isEqualTo(strategySpy);
verifyNoInteractions(strategySpy);
}
@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;
long lastReportedTime = 312L;
ActivityStrategy strategySpy = spy(new FirstEventActivityStrategy());
ActivityState<TransportProtos.SessionInfoProto> state = new ActivityState<>();
state.setLastRecordedTime(lastRecordedTime);
state.setLastReportedTime(lastReportedTime);
state.setMetadata(TransportProtos.SessionInfoProto.getDefaultInstance());
state.setStrategy(strategySpy);
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.getLastReportedTime()).isEqualTo(lastReportedTime);
assertThat(updatedState.getMetadata()).isEqualTo(sessionInfo);
assertThat(updatedState.getStrategy()).isEqualTo(strategySpy);
verifyNoInteractions(strategySpy);
}
@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