diff --git a/application/src/test/java/org/thingsboard/server/service/integration/IntegrationActivityManagerTest.java b/application/src/test/java/org/thingsboard/server/service/integration/IntegrationActivityManagerTest.java new file mode 100644 index 0000000000..6748bced4d --- /dev/null +++ b/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> 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 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> msgCaptor = ArgumentCaptor.forClass(TbProtoQueueMsg.class); + ArgumentCaptor callbackCaptor = ArgumentCaptor.forClass(TbQueueCallback.class); + verify(tbCoreMsgProducerMock).send(eq(tpi), msgCaptor.capture(), callbackCaptor.capture()); + + TbProtoQueueMsg 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 reportingStrategyClass + ) { + // GIVEN + var key = new IntegrationActivityKey(TENANT_ID, DEVICE_ID); + when(integrationServiceMock.createNewState(key)).thenCallRealMethod(); + ReflectionTestUtils.setField(integrationServiceMock, "reportingStrategyName", reportingStrategyName); + + // WHEN + ActivityState 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 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 expectedState = new ActivityState<>(); + expectedState.setLastRecordedTime(expectedLastRecordedTime); + expectedState.setLastReportedTime(expectedLastReportedTime); + expectedState.setStrategy(expectedStrategy); + + when(integrationServiceMock.updateState(key, expectedState)).thenCallRealMethod(); + + // WHEN + ActivityState 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 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 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)); + } + +} diff --git a/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/activity/AbstractActivityManager.java b/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/activity/AbstractActivityManager.java index 1090c7c341..b5f4bf8f9e 100644 --- a/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/activity/AbstractActivityManager.java +++ b/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/activity/AbstractActivityManager.java @@ -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 implements Activity protected abstract ActivityState createNewState(Key key); - protected abstract ActivityStrategy getStrategy(); - protected abstract ActivityState updateState(Key key, ActivityState state); - protected abstract boolean hasExpired(Key key, ActivityState 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 callback); @@ -102,7 +99,6 @@ public abstract class AbstractActivityManager implements Activity return null; } state = newState; - state.setStrategy(getStrategy()); } if (state.getLastRecordedTime() < newLastRecordedTime) { state.setLastRecordedTime(newLastRecordedTime); @@ -156,7 +152,7 @@ public abstract class AbstractActivityManager 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 implements Activity if (hasExpired) { states.remove(key); - onStateExpire(key, metadata); + onStateExpiry(key, metadata); shouldReport = true; } diff --git a/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/service/DefaultTransportService.java b/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/service/DefaultTransportService.java index 4b70638e37..9cd41e7a62 100644 --- a/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/service/DefaultTransportService.java +++ b/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 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 updateState(UUID sessionId, ActivityState state) { SessionMetaData session = sessions.get(sessionId); @@ -835,12 +830,12 @@ public class DefaultTransportService extends AbstractActivityManager 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) {