Browse Source

Minor refactoring and tests for integration activity manager

pull/9980/head
Dmytro Skarzhynets 3 years ago
parent
commit
76fec46720
  1. 299
      application/src/test/java/org/thingsboard/server/service/integration/IntegrationActivityManagerTest.java
  2. 12
      common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/activity/AbstractActivityManager.java
  3. 13
      common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/service/DefaultTransportService.java

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

@ -0,0 +1,299 @@
/**
* ThingsBoard, Inc. ("COMPANY") CONFIDENTIAL
*
* Copyright © 2016-2023 ThingsBoard, Inc. All Rights Reserved.
*
* NOTICE: All information contained herein is, and remains
* the property of ThingsBoard, Inc. and its suppliers,
* if any. The intellectual and technical concepts contained
* herein are proprietary to ThingsBoard, Inc.
* and its suppliers and may be covered by U.S. and Foreign Patents,
* patents in process, and are protected by trade secret or copyright law.
*
* Dissemination of this information or reproduction of this material is strictly forbidden
* unless prior written permission is obtained from COMPANY.
*
* Access to the source code contained herein is hereby forbidden to anyone except current COMPANY employees,
* managers or contractors who have executed Confidentiality and Non-disclosure agreements
* explicitly covering such access.
*
* The copyright notice above does not evidence any actual or intended publication
* or disclosure of this source code, which includes
* information that is confidential and/or proprietary, and is a trade secret, of COMPANY.
* ANY REPRODUCTION, MODIFICATION, DISTRIBUTION, PUBLIC PERFORMANCE,
* OR PUBLIC DISPLAY OF OR THROUGH USE OF THIS SOURCE CODE WITHOUT
* THE EXPRESS WRITTEN CONSENT OF COMPANY IS STRICTLY PROHIBITED,
* AND IN VIOLATION OF APPLICABLE LAWS AND INTERNATIONAL TREATIES.
* THE RECEIPT OR POSSESSION OF THIS SOURCE CODE AND/OR RELATED INFORMATION
* DOES NOT CONVEY OR IMPLY ANY RIGHTS TO REPRODUCE, DISCLOSE OR DISTRIBUTE ITS CONTENTS,
* OR TO MANUFACTURE, USE, OR SELL ANYTHING THAT IT MAY DESCRIBE, IN WHOLE OR IN PART.
*/
package org.thingsboard.server.service.integration;
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.data.id.DeviceId;
import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.common.msg.queue.ServiceType;
import org.thingsboard.server.common.msg.queue.TopicPartitionInfo;
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 org.thingsboard.server.queue.TbQueueCallback;
import org.thingsboard.server.queue.TbQueueProducer;
import org.thingsboard.server.queue.common.TbProtoQueueMsg;
import org.thingsboard.server.queue.discovery.PartitionService;
import java.util.UUID;
import java.util.stream.Stream;
import static org.assertj.core.api.Assertions.assertThat;
import static org.assertj.core.api.Assertions.assertThatNoException;
import static org.mockito.ArgumentMatchers.any;
import static org.mockito.ArgumentMatchers.eq;
import static org.mockito.Mockito.doCallRealMethod;
import static org.mockito.Mockito.doNothing;
import static org.mockito.Mockito.mock;
import static org.mockito.Mockito.spy;
import static org.mockito.Mockito.verify;
import static org.mockito.Mockito.verifyNoInteractions;
import static org.mockito.Mockito.when;
@ExtendWith(MockitoExtension.class)
public class IntegrationActivityManagerTest {
private final TenantId TENANT_ID = TenantId.fromUUID(UUID.fromString("1306648a-9b26-11ee-b9d1-0242ac120002"));
private final DeviceId DEVICE_ID = DeviceId.fromString("1d288a06-9b26-11ee-b9d1-0242ac120002");
@Mock
private DefaultPlatformIntegrationService integrationServiceMock;
@Test
void testReportActivity() {
// GIVEN
var key = new IntegrationActivityKey(TENANT_ID, DEVICE_ID);
TbQueueProducer<TbProtoQueueMsg<TransportProtos.ToCoreMsg>> tbCoreMsgProducerMock = mock(TbQueueProducer.class);
ReflectionTestUtils.setField(integrationServiceMock, "tbCoreMsgProducer", tbCoreMsgProducerMock);
PartitionService partitionServiceMock = mock(PartitionService.class);
ReflectionTestUtils.setField(integrationServiceMock, "partitionService", partitionServiceMock);
TopicPartitionInfo tpi = TopicPartitionInfo.builder().build();
when(partitionServiceMock.resolve(ServiceType.TB_CORE, TENANT_ID, DEVICE_ID)).thenReturn(tpi);
ActivityReportCallback<IntegrationActivityKey> callbackMock = mock(ActivityReportCallback.class);
long expectedTime = 123L;
doCallRealMethod().when(integrationServiceMock).reportActivity(key, null, expectedTime, callbackMock);
// WHEN
integrationServiceMock.reportActivity(key, null, expectedTime, callbackMock);
// THEN
verify(partitionServiceMock).resolve(ServiceType.TB_CORE, TENANT_ID, DEVICE_ID);
ArgumentCaptor<TbProtoQueueMsg<TransportProtos.ToCoreMsg>> msgCaptor = ArgumentCaptor.forClass(TbProtoQueueMsg.class);
ArgumentCaptor<TbQueueCallback> callbackCaptor = ArgumentCaptor.forClass(TbQueueCallback.class);
verify(tbCoreMsgProducerMock).send(eq(tpi), msgCaptor.capture(), callbackCaptor.capture());
TbProtoQueueMsg<TransportProtos.ToCoreMsg> queueMsg = msgCaptor.getValue();
TransportProtos.DeviceActivityProto expectedDeviceActivityMsg = TransportProtos.DeviceActivityProto.newBuilder()
.setTenantIdMSB(TENANT_ID.getId().getMostSignificantBits())
.setTenantIdLSB(TENANT_ID.getId().getLeastSignificantBits())
.setDeviceIdMSB(DEVICE_ID.getId().getMostSignificantBits())
.setDeviceIdLSB(DEVICE_ID.getId().getLeastSignificantBits())
.setLastActivityTime(expectedTime)
.build();
TransportProtos.ToCoreMsg expectedToCoreMsg = TransportProtos.ToCoreMsg.newBuilder()
.setDeviceActivityMsg(expectedDeviceActivityMsg)
.build();
assertThat(queueMsg.getKey()).isEqualTo(DEVICE_ID.getId());
assertThat(queueMsg.getValue()).isEqualTo(expectedToCoreMsg);
TbQueueCallback queueCallback = callbackCaptor.getValue();
queueCallback.onSuccess(null);
verify(callbackMock).onSuccess(key, expectedTime);
var throwable = new Throwable();
queueCallback.onFailure(throwable);
verify(callbackMock).onFailure(key, throwable);
}
@Test
void givenPostTelemetryMsg_whenProcessingMsg_thenShouldCallOnActivity() {
// GIVEN
var key = new IntegrationActivityKey(TENANT_ID, DEVICE_ID);
TransportProtos.SessionInfoProto sessionInfo = TransportProtos.SessionInfoProto.newBuilder()
.setTenantIdMSB(TENANT_ID.getId().getMostSignificantBits())
.setTenantIdLSB(TENANT_ID.getId().getLeastSignificantBits())
.setDeviceIdMSB(DEVICE_ID.getId().getMostSignificantBits())
.setDeviceIdLSB(DEVICE_ID.getId().getLeastSignificantBits())
.setDeviceName("Test Device")
.setDeviceType("default")
.build();
TransportProtos.PostTelemetryMsg postTelemetryMsg = TransportProtos.PostTelemetryMsg.getDefaultInstance();
doCallRealMethod().when(integrationServiceMock).process(sessionInfo, postTelemetryMsg, null);
// WHEN
integrationServiceMock.process(sessionInfo, postTelemetryMsg, null);
// THEN
verify(integrationServiceMock).onActivity(key);
}
@Test
void givenPostAttributesMsg_whenProcessingMsg_thenShouldCallOnActivity() {
// GIVEN
var key = new IntegrationActivityKey(TENANT_ID, DEVICE_ID);
TransportProtos.SessionInfoProto sessionInfo = TransportProtos.SessionInfoProto.newBuilder()
.setTenantIdMSB(TENANT_ID.getId().getMostSignificantBits())
.setTenantIdLSB(TENANT_ID.getId().getLeastSignificantBits())
.setDeviceIdMSB(DEVICE_ID.getId().getMostSignificantBits())
.setDeviceIdLSB(DEVICE_ID.getId().getLeastSignificantBits())
.setDeviceName("Test Device")
.setDeviceType("default")
.build();
TransportProtos.PostAttributeMsg postAttributeMsg = TransportProtos.PostAttributeMsg.getDefaultInstance();
doCallRealMethod().when(integrationServiceMock).process(sessionInfo, postAttributeMsg, null);
doNothing().when(integrationServiceMock).sendToRuleEngine(any(), any(), any(), any(), any(), any(), any());
// WHEN
integrationServiceMock.process(sessionInfo, postAttributeMsg, null);
// THEN
verify(integrationServiceMock).onActivity(key);
}
@ParameterizedTest
@MethodSource("provideTestParamsForCreateNewState")
void givenDifferentReportingStrategies_whenCreatingNewState_thenShouldCreateEmptyStateWithCorrectStrategy(
String reportingStrategyName, Class<ActivityStrategy> reportingStrategyClass
) {
// GIVEN
var key = new IntegrationActivityKey(TENANT_ID, DEVICE_ID);
when(integrationServiceMock.createNewState(key)).thenCallRealMethod();
ReflectionTestUtils.setField(integrationServiceMock, "reportingStrategyName", reportingStrategyName);
// WHEN
ActivityState<Void> newState = integrationServiceMock.createNewState(key);
// THEN
assertThat(newState).isNotNull();
assertThat(newState.getLastRecordedTime()).isEqualTo(0L);
assertThat(newState.getLastReportedTime()).isEqualTo(0L);
assertThat(newState.getStrategy()).isInstanceOf(reportingStrategyClass);
}
private static Stream<Arguments> provideTestParamsForCreateNewState() {
return Stream.of(
Arguments.of("ALL", AllEventsActivityStrategy.class),
Arguments.of("FIRST", FirstEventActivityStrategy.class),
Arguments.of("LAST", LastEventActivityStrategy.class),
Arguments.of("FIRST_AND_LAST", FirstAndLastEventActivityStrategy.class)
);
}
@Test
void givenActivityState_whenUpdatingActivityState_thenShouldReturnSameInstanceWithNoChanges() {
// GIVEN
var key = new IntegrationActivityKey(TENANT_ID, DEVICE_ID);
long expectedLastRecordedTime = 123L;
long expectedLastReportedTime = 312L;
ActivityStrategy expectedStrategy = spy(new FirstEventActivityStrategy());
ActivityState<Void> expectedState = new ActivityState<>();
expectedState.setLastRecordedTime(expectedLastRecordedTime);
expectedState.setLastReportedTime(expectedLastReportedTime);
expectedState.setStrategy(expectedStrategy);
when(integrationServiceMock.updateState(key, expectedState)).thenCallRealMethod();
// WHEN
ActivityState<Void> actualNewState = integrationServiceMock.updateState(key, expectedState);
// THEN
assertThat(actualNewState).isSameAs(expectedState);
assertThat(actualNewState.getLastRecordedTime()).isEqualTo(expectedLastRecordedTime);
assertThat(actualNewState.getLastReportedTime()).isEqualTo(expectedLastReportedTime);
assertThat(actualNewState.getStrategy()).isSameAs(expectedStrategy);
verifyNoInteractions(expectedStrategy);
assertThat(actualNewState.getMetadata()).isNull();
}
@ParameterizedTest
@MethodSource("provideTestParamsForHasExpiredTrue")
public void givenExpiredLastRecordedTime_whenCheckingForExpiry_thenShouldReturnTrue(long currentTimeMillis, long lastRecordedTime, long reportingPeriodMillis) {
// GIVEN
ReflectionTestUtils.setField(integrationServiceMock, "reportingPeriodMillis", reportingPeriodMillis);
when(integrationServiceMock.getCurrentTimeMillis()).thenReturn(currentTimeMillis);
when(integrationServiceMock.hasExpired(lastRecordedTime)).thenCallRealMethod();
// WHEN
boolean hasExpired = integrationServiceMock.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)
);
}
@ParameterizedTest
@MethodSource("provideTestParamsForHasExpiredFalse")
public void givenNotExpiredLastRecordedTime_whenCheckingForExpiry_thenShouldReturnFalse(long currentTimeMillis, long lastRecordedTime, long reportingPeriodMillis) {
// GIVEN
ReflectionTestUtils.setField(integrationServiceMock, "reportingPeriodMillis", reportingPeriodMillis);
when(integrationServiceMock.getCurrentTimeMillis()).thenReturn(currentTimeMillis);
when(integrationServiceMock.hasExpired(lastRecordedTime)).thenCallRealMethod();
// WHEN
boolean hasExpired = integrationServiceMock.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)
);
}
@Test
void givenKeyAndVoidMetadata_whenOnStateExpiryCalled_thenShouldDoNothing() {
// GIVEN
var key = new IntegrationActivityKey(TENANT_ID, DEVICE_ID);
doCallRealMethod().when(integrationServiceMock).onStateExpiry(key, null);
// WHEN-THEN
assertThatNoException().isThrownBy(() -> integrationServiceMock.onStateExpiry(key, null));
}
}

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

@ -33,7 +33,6 @@ package org.thingsboard.server.common.transport.activity;
import lombok.extern.slf4j.Slf4j;
import org.springframework.beans.factory.annotation.Autowired;
import org.thingsboard.server.common.data.StringUtils;
import org.thingsboard.server.common.transport.activity.strategy.ActivityStrategy;
import org.thingsboard.server.queue.scheduler.SchedulerComponent;
import java.util.Map;
@ -71,13 +70,11 @@ public abstract class AbstractActivityManager<Key, Metadata> implements Activity
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(Key key, ActivityState<Metadata> state);
protected abstract boolean hasExpired(long lastRecordedTime);
protected abstract void onStateExpire(Key key, Metadata metadata);
protected abstract void onStateExpiry(Key key, Metadata metadata);
protected abstract void reportActivity(Key key, Metadata metadata, long timeToReport, ActivityReportCallback<Key> callback);
@ -102,7 +99,6 @@ public abstract class AbstractActivityManager<Key, Metadata> implements Activity
return null;
}
state = newState;
state.setStrategy(getStrategy());
}
if (state.getLastRecordedTime() < newLastRecordedTime) {
state.setLastRecordedTime(newLastRecordedTime);
@ -156,7 +152,7 @@ public abstract class AbstractActivityManager<Key, Metadata> implements Activity
lastRecordedTime = updatedState.getLastRecordedTime();
lastReportedTime = updatedState.getLastReportedTime();
metadata = updatedState.getMetadata();
hasExpired = hasExpired(key, updatedState);
hasExpired = hasExpired(lastRecordedTime);
shouldReport = updatedState.getStrategy().onReportingPeriodEnd();
} else {
states.remove(key);
@ -166,7 +162,7 @@ public abstract class AbstractActivityManager<Key, Metadata> implements Activity
if (hasExpired) {
states.remove(key);
onStateExpire(key, metadata);
onStateExpiry(key, metadata);
shouldReport = true;
}

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

@ -74,7 +74,6 @@ import org.thingsboard.server.common.transport.TransportTenantProfileCache;
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.ActivityStrategyFactory;
import org.thingsboard.server.common.transport.auth.GetOrCreateDeviceFromGatewayResponse;
import org.thingsboard.server.common.transport.auth.TransportDeviceInfo;
@ -797,14 +796,10 @@ public class DefaultTransportService extends AbstractActivityManager<UUID, Trans
}
ActivityState<TransportProtos.SessionInfoProto> state = new ActivityState<>();
state.setMetadata(session.getSessionInfo());
state.setStrategy(ActivityStrategyFactory.createStrategy(reportingStrategyName));
return state;
}
@Override
protected ActivityStrategy getStrategy() {
return ActivityStrategyFactory.createStrategy(reportingStrategyName);
}
@Override
protected ActivityState<TransportProtos.SessionInfoProto> updateState(UUID sessionId, ActivityState<TransportProtos.SessionInfoProto> state) {
SessionMetaData session = sessions.get(sessionId);
@ -835,12 +830,12 @@ public class DefaultTransportService extends AbstractActivityManager<UUID, Trans
}
@Override
protected boolean hasExpired(UUID uuid, ActivityState<TransportProtos.SessionInfoProto> state) {
return (System.currentTimeMillis() - sessionInactivityTimeout) > state.getLastRecordedTime();
protected boolean hasExpired(long lastRecordedTime) {
return (System.currentTimeMillis() - sessionInactivityTimeout) > lastRecordedTime;
}
@Override
protected void onStateExpire(UUID sessionId, TransportProtos.SessionInfoProto sessionInfo) {
protected void onStateExpiry(UUID sessionId, TransportProtos.SessionInfoProto sessionInfo) {
log.debug("[{}] Session with id: [{}] has expired due to last activity time.", name, sessionId);
SessionMetaData expiredSession = sessions.remove(sessionId);
if (expiredSession != null) {

Loading…
Cancel
Save