From 8a6484203972f2c6b70021fb896680fbe5dda95e Mon Sep 17 00:00:00 2001 From: Dmytro Skarzhynets Date: Wed, 30 Aug 2023 10:44:39 +0300 Subject: [PATCH 1/9] Add missing sleep, fix wrong timeout expiry check --- .../server/service/state/DefaultDeviceStateService.java | 2 +- .../server/service/state/DefaultDeviceStateServiceTest.java | 1 + 2 files changed, 2 insertions(+), 1 deletion(-) diff --git a/application/src/main/java/org/thingsboard/server/service/state/DefaultDeviceStateService.java b/application/src/main/java/org/thingsboard/server/service/state/DefaultDeviceStateService.java index cdc5afe361..f9afb5b2e5 100644 --- a/application/src/main/java/org/thingsboard/server/service/state/DefaultDeviceStateService.java +++ b/application/src/main/java/org/thingsboard/server/service/state/DefaultDeviceStateService.java @@ -486,7 +486,7 @@ public class DefaultDeviceStateService extends AbstractPartitionBasedService Date: Wed, 30 Aug 2023 14:03:43 +0300 Subject: [PATCH 2/9] Add missing sleep in `decreaseInactivityForInactiveDeviceTest()` --- .../server/service/state/DefaultDeviceStateServiceTest.java | 1 + 1 file changed, 1 insertion(+) diff --git a/application/src/test/java/org/thingsboard/server/service/state/DefaultDeviceStateServiceTest.java b/application/src/test/java/org/thingsboard/server/service/state/DefaultDeviceStateServiceTest.java index c0690c56e1..ceb98480c1 100644 --- a/application/src/test/java/org/thingsboard/server/service/state/DefaultDeviceStateServiceTest.java +++ b/application/src/test/java/org/thingsboard/server/service/state/DefaultDeviceStateServiceTest.java @@ -330,6 +330,7 @@ public class DefaultDeviceStateServiceTest { Mockito.reset(telemetrySubscriptionService); long newTimeout = 1; + Thread.sleep(newTimeout); service.onDeviceInactivityTimeoutUpdate(tenantId, deviceId, newTimeout); Mockito.verify(telemetrySubscriptionService, Mockito.never()).saveAttrAndNotify(Mockito.any(), Mockito.eq(deviceId), Mockito.any(), Mockito.eq("active"), Mockito.any(), Mockito.any()); From 7d13b034a279c53d30c6b1b8867913b6cfc76b99 Mon Sep 17 00:00:00 2001 From: Dmytro Skarzhynets Date: Sun, 3 Sep 2023 11:22:30 +0300 Subject: [PATCH 3/9] Migrate test to Junit 5, perform small cleanup --- .../state/DefaultDeviceStateServiceTest.java | 83 ++++++++++--------- 1 file changed, 42 insertions(+), 41 deletions(-) diff --git a/application/src/test/java/org/thingsboard/server/service/state/DefaultDeviceStateServiceTest.java b/application/src/test/java/org/thingsboard/server/service/state/DefaultDeviceStateServiceTest.java index ceb98480c1..e03e4e94cf 100644 --- a/application/src/test/java/org/thingsboard/server/service/state/DefaultDeviceStateServiceTest.java +++ b/application/src/test/java/org/thingsboard/server/service/state/DefaultDeviceStateServiceTest.java @@ -15,13 +15,11 @@ */ package org.thingsboard.server.service.state; -import org.junit.Assert; -import org.junit.Before; -import org.junit.Test; -import org.junit.runner.RunWith; +import org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.api.Test; +import org.junit.jupiter.api.extension.ExtendWith; import org.mockito.Mock; -import org.mockito.Mockito; -import org.mockito.junit.MockitoJUnitRunner; +import org.mockito.junit.jupiter.MockitoExtension; import org.springframework.test.util.ReflectionTestUtils; import org.thingsboard.server.cluster.TbClusterService; import org.thingsboard.server.common.data.DeviceIdInfo; @@ -49,16 +47,20 @@ import java.util.List; import java.util.Map; import java.util.UUID; -import static org.hamcrest.CoreMatchers.is; -import static org.hamcrest.MatcherAssert.assertThat; +import static org.assertj.core.api.Assertions.assertThat; +import static org.mockito.ArgumentMatchers.any; +import static org.mockito.ArgumentMatchers.eq; import static org.mockito.BDDMockito.willReturn; -import static org.mockito.Mockito.mock; import static org.mockito.Mockito.never; +import static org.mockito.Mockito.reset; import static org.mockito.Mockito.spy; import static org.mockito.Mockito.times; +import static org.mockito.Mockito.verify; +import static org.mockito.Mockito.when; +import static org.thingsboard.server.service.state.DefaultDeviceStateService.ACTIVITY_STATE; import static org.thingsboard.server.service.state.DefaultDeviceStateService.INACTIVITY_TIMEOUT; -@RunWith(MockitoJUnitRunner.class) +@ExtendWith(MockitoExtension.class) public class DefaultDeviceStateServiceTest { @Mock @@ -75,6 +77,10 @@ public class DefaultDeviceStateServiceTest { DeviceStateData deviceStateDataMock; @Mock EntityQueryRepository entityQueryRepository; + @Mock + TelemetrySubscriptionService telemetrySubscriptionService; + @Mock + NotificationRuleProcessor notificationRuleProcessor; TenantId tenantId = new TenantId(UUID.fromString("00797a3b-7aeb-4b5b-b57a-c2a810d0f112")); DeviceId deviceId = DeviceId.fromString("00797a3b-7aeb-4b5b-b57a-c2a810d0f112"); @@ -82,31 +88,23 @@ public class DefaultDeviceStateServiceTest { DefaultDeviceStateService service; - TelemetrySubscriptionService telemetrySubscriptionService; - - @Before + @BeforeEach public void setUp() { - service = spy(new DefaultDeviceStateService(deviceService, attributesService, tsService, clusterService, partitionService, entityQueryRepository, null, null, mock(NotificationRuleProcessor.class))); - telemetrySubscriptionService = Mockito.mock(TelemetrySubscriptionService.class); + service = spy(new DefaultDeviceStateService(deviceService, attributesService, tsService, clusterService, partitionService, entityQueryRepository, null, null, notificationRuleProcessor)); ReflectionTestUtils.setField(service, "tsSubService", telemetrySubscriptionService); ReflectionTestUtils.setField(service, "defaultStateCheckIntervalInSec", 60); ReflectionTestUtils.setField(service, "defaultActivityStatsIntervalInSec", 60); ReflectionTestUtils.setField(service, "initFetchPackSize", 10); tpi = TopicPartitionInfo.builder().myPartition(true).build(); - Mockito.when(partitionService.resolve(ServiceType.TB_CORE, tenantId, deviceId)).thenReturn(tpi); - Mockito.when(entityQueryRepository.findEntityDataByQueryInternal(Mockito.any())).thenReturn(new PageData<>()); - var deviceIdInfo = new DeviceIdInfo(tenantId.getId(), null, deviceId.getId()); - Mockito.when(deviceService.findDeviceIdInfos(Mockito.any())) - .thenReturn(new PageData<>(List.of(deviceIdInfo), 0, 1, false)); } @Test public void givenDeviceIdFromDeviceStatesMap_whenGetOrFetchDeviceStateData_thenNoStackOverflow() { service.deviceStates.put(deviceId, deviceStateDataMock); DeviceStateData deviceStateData = service.getOrFetchDeviceStateData(deviceId); - assertThat(deviceStateData, is(deviceStateDataMock)); - Mockito.verify(service, never()).fetchDeviceStateDataUsingEntityDataQuery(deviceId); + assertThat(deviceStateData).isEqualTo(deviceStateDataMock); + verify(service, never()).fetchDeviceStateDataUsingEntityDataQuery(deviceId); } @Test @@ -114,8 +112,8 @@ public class DefaultDeviceStateServiceTest { service.deviceStates.clear(); willReturn(deviceStateDataMock).given(service).fetchDeviceStateDataUsingEntityDataQuery(deviceId); DeviceStateData deviceStateData = service.getOrFetchDeviceStateData(deviceId); - assertThat(deviceStateData, is(deviceStateDataMock)); - Mockito.verify(service, times(1)).fetchDeviceStateDataUsingEntityDataQuery(deviceId); + assertThat(deviceStateData).isEqualTo(deviceStateDataMock); + verify(service, times(1)).fetchDeviceStateDataUsingEntityDataQuery(deviceId); } @Test @@ -151,14 +149,19 @@ public class DefaultDeviceStateServiceTest { DeviceStateData deviceStateData = service.toDeviceStateData(new EntityData(deviceId, latest, Map.of()), new DeviceIdInfo(TenantId.SYS_TENANT_ID.getId(), UUID.randomUUID(), deviceUuid)); - Assert.assertEquals(5000L, deviceStateData.getState().getInactivityTimeout()); + assertThat(deviceStateData.getState().getInactivityTimeout()).isEqualTo(5000L); } private void initStateService(long timeout) throws InterruptedException { service.stop(); - Mockito.reset(service, telemetrySubscriptionService); - ReflectionTestUtils.setField(service, "defaultInactivityTimeoutMs", timeout); + reset(service, telemetrySubscriptionService); + service.setDefaultInactivityTimeoutMs(timeout); service.init(); + when(partitionService.resolve(ServiceType.TB_CORE, tenantId, deviceId)).thenReturn(tpi); + when(entityQueryRepository.findEntityDataByQueryInternal(any())).thenReturn(new PageData<>()); + var deviceIdInfo = new DeviceIdInfo(tenantId.getId(), null, deviceId.getId()); + when(deviceService.findDeviceIdInfos(any())) + .thenReturn(new PageData<>(List.of(deviceIdInfo), 0, 1, false)); PartitionChangeEvent event = new PartitionChangeEvent(this, new QueueKey(ServiceType.TB_CORE), Collections.singleton(tpi)); service.onApplicationEvent(event); Thread.sleep(100); @@ -185,7 +188,7 @@ public class DefaultDeviceStateServiceTest { service.checkStates(); activityVerify(false); - Mockito.reset(telemetrySubscriptionService); + reset(telemetrySubscriptionService); long increase = 100; long newTimeout = System.currentTimeMillis() - deviceState.getLastActivityTime() + increase; @@ -196,7 +199,7 @@ public class DefaultDeviceStateServiceTest { service.checkStates(); activityVerify(false); - Mockito.reset(telemetrySubscriptionService); + reset(telemetrySubscriptionService); service.onDeviceActivity(tenantId, deviceId, System.currentTimeMillis()); activityVerify(true); @@ -223,18 +226,18 @@ public class DefaultDeviceStateServiceTest { service.onDeviceActivity(tenantId, deviceId, System.currentTimeMillis()); activityVerify(true); - Mockito.reset(telemetrySubscriptionService); + reset(telemetrySubscriptionService); long increase = 100; long newTimeout = System.currentTimeMillis() - deviceState.getLastActivityTime() + increase; service.onDeviceInactivityTimeoutUpdate(tenantId, deviceId, newTimeout); - Mockito.verify(telemetrySubscriptionService, Mockito.never()).saveAttrAndNotify(Mockito.any(), Mockito.eq(deviceId), Mockito.any(), Mockito.eq("active"), Mockito.any(), Mockito.any()); + verify(telemetrySubscriptionService, never()).saveAttrAndNotify(any(), eq(deviceId), any(), eq(ACTIVITY_STATE), any(), any()); Thread.sleep(defaultTimeout + increase); service.checkStates(); activityVerify(false); - Mockito.reset(telemetrySubscriptionService); + reset(telemetrySubscriptionService); service.onDeviceActivity(tenantId, deviceId, System.currentTimeMillis()); activityVerify(true); @@ -264,11 +267,11 @@ public class DefaultDeviceStateServiceTest { service.checkStates(); activityVerify(false); - Mockito.reset(telemetrySubscriptionService); + reset(telemetrySubscriptionService); long newTimeout = 1; Thread.sleep(newTimeout); - Mockito.verify(telemetrySubscriptionService, Mockito.never()).saveAttrAndNotify(Mockito.any(), Mockito.eq(deviceId), Mockito.any(), Mockito.eq("active"), Mockito.any(), Mockito.any()); + verify(telemetrySubscriptionService, never()).saveAttrAndNotify(any(), eq(deviceId), any(), eq(ACTIVITY_STATE), any(), any()); } @Test @@ -289,16 +292,14 @@ public class DefaultDeviceStateServiceTest { service.onDeviceActivity(tenantId, deviceId, System.currentTimeMillis()); activityVerify(true); - Mockito.reset(telemetrySubscriptionService); - - Mockito.verify(telemetrySubscriptionService, Mockito.never()).saveAttrAndNotify(Mockito.any(), Mockito.eq(deviceId), Mockito.any(), Mockito.eq("active"), Mockito.any(), Mockito.any()); + verify(telemetrySubscriptionService, never()).saveAttrAndNotify(any(), eq(deviceId), any(), eq(ACTIVITY_STATE), any(), any()); long newTimeout = 1; Thread.sleep(newTimeout); service.onDeviceInactivityTimeoutUpdate(tenantId, deviceId, newTimeout); activityVerify(false); - Mockito.reset(telemetrySubscriptionService); + reset(telemetrySubscriptionService); service.onDeviceInactivityTimeoutUpdate(tenantId, deviceId, defaultTimeout); activityVerify(true); @@ -327,17 +328,17 @@ public class DefaultDeviceStateServiceTest { Thread.sleep(defaultTimeout); service.checkStates(); activityVerify(false); - Mockito.reset(telemetrySubscriptionService); + reset(telemetrySubscriptionService); long newTimeout = 1; Thread.sleep(newTimeout); service.onDeviceInactivityTimeoutUpdate(tenantId, deviceId, newTimeout); - Mockito.verify(telemetrySubscriptionService, Mockito.never()).saveAttrAndNotify(Mockito.any(), Mockito.eq(deviceId), Mockito.any(), Mockito.eq("active"), Mockito.any(), Mockito.any()); + verify(telemetrySubscriptionService, never()).saveAttrAndNotify(any(), eq(deviceId), any(), eq(ACTIVITY_STATE), any(), any()); } private void activityVerify(boolean isActive) { - Mockito.verify(telemetrySubscriptionService, Mockito.times(1)).saveAttrAndNotify(Mockito.any(), Mockito.eq(deviceId), Mockito.any(), Mockito.eq("active"), Mockito.eq(isActive), Mockito.any()); + verify(telemetrySubscriptionService, times(1)).saveAttrAndNotify(any(), eq(deviceId), any(), eq(ACTIVITY_STATE), eq(isActive), any()); } } \ No newline at end of file From 4180866cccd491917b3eb1779cdc6bc8b71ae616 Mon Sep 17 00:00:00 2001 From: Dmytro Skarzhynets Date: Wed, 6 Sep 2023 18:37:42 +0300 Subject: [PATCH 4/9] Add parameterized test to cover `updateInactivityStateIfExpired()` --- .../state/DefaultDeviceStateServiceTest.java | 183 +++++++++++++++++- 1 file changed, 182 insertions(+), 1 deletion(-) diff --git a/application/src/test/java/org/thingsboard/server/service/state/DefaultDeviceStateServiceTest.java b/application/src/test/java/org/thingsboard/server/service/state/DefaultDeviceStateServiceTest.java index e03e4e94cf..10c507f884 100644 --- a/application/src/test/java/org/thingsboard/server/service/state/DefaultDeviceStateServiceTest.java +++ b/application/src/test/java/org/thingsboard/server/service/state/DefaultDeviceStateServiceTest.java @@ -18,6 +18,10 @@ package org.thingsboard.server.service.state; 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; @@ -25,10 +29,13 @@ import org.thingsboard.server.cluster.TbClusterService; import org.thingsboard.server.common.data.DeviceIdInfo; import org.thingsboard.server.common.data.id.DeviceId; import org.thingsboard.server.common.data.id.TenantId; +import org.thingsboard.server.common.data.msg.TbMsgType; +import org.thingsboard.server.common.data.notification.rule.trigger.DeviceActivityTrigger; import org.thingsboard.server.common.data.page.PageData; import org.thingsboard.server.common.data.query.EntityData; import org.thingsboard.server.common.data.query.EntityKeyType; import org.thingsboard.server.common.data.query.TsValue; +import org.thingsboard.server.common.msg.TbMsg; import org.thingsboard.server.common.msg.TbMsgMetaData; import org.thingsboard.server.common.msg.notification.NotificationRuleProcessor; import org.thingsboard.server.common.msg.queue.ServiceType; @@ -46,10 +53,13 @@ import java.util.Collections; import java.util.List; import java.util.Map; import java.util.UUID; +import java.util.stream.Stream; import static org.assertj.core.api.Assertions.assertThat; import static org.mockito.ArgumentMatchers.any; import static org.mockito.ArgumentMatchers.eq; +import static org.mockito.BDDMockito.given; +import static org.mockito.BDDMockito.then; import static org.mockito.BDDMockito.willReturn; import static org.mockito.Mockito.never; import static org.mockito.Mockito.reset; @@ -57,7 +67,9 @@ import static org.mockito.Mockito.spy; import static org.mockito.Mockito.times; import static org.mockito.Mockito.verify; import static org.mockito.Mockito.when; +import static org.thingsboard.server.common.data.DataConstants.SERVER_SCOPE; import static org.thingsboard.server.service.state.DefaultDeviceStateService.ACTIVITY_STATE; +import static org.thingsboard.server.service.state.DefaultDeviceStateService.INACTIVITY_ALARM_TIME; import static org.thingsboard.server.service.state.DefaultDeviceStateService.INACTIVITY_TIMEOUT; @ExtendWith(MockitoExtension.class) @@ -341,4 +353,173 @@ public class DefaultDeviceStateServiceTest { verify(telemetrySubscriptionService, times(1)).saveAttrAndNotify(any(), eq(deviceId), any(), eq(ACTIVITY_STATE), eq(isActive), any()); } -} \ No newline at end of file + @Test + public void givenStateDataIsNull_whenUpdateInactivityTimeoutIfExpired_thenShouldCleanupDevice() { + // GIVEN + service.deviceStates.put(deviceId, deviceStateDataMock); + + // WHEN + service.updateInactivityStateIfExpired(System.currentTimeMillis(), deviceId, null); + + // THEN + assertThat(service.deviceStates.get(deviceId)).isNull(); + assertThat(service.deviceStates.size()).isEqualTo(0); + assertThat(service.deviceStates.isEmpty()).isTrue(); + } + + @Test + public void givenNotMyPartition_whenUpdateInactivityTimeoutIfExpired_thenShouldCleanupDevice() { + // GIVEN + long currentTime = System.currentTimeMillis(); + + DeviceState deviceState = DeviceState.builder() + .active(true) + .lastConnectTime(currentTime - 8000) + .lastActivityTime(currentTime - 4000) + .lastDisconnectTime(0) + .lastInactivityAlarmTime(0) + .inactivityTimeout(3000) + .build(); + + DeviceStateData stateData = DeviceStateData.builder() + .tenantId(tenantId) + .deviceId(deviceId) + .deviceCreationTime(currentTime - 10000) + .state(deviceState) + .build(); + + service.deviceStates.put(deviceId, stateData); + + var notMyTpi = TopicPartitionInfo.builder().myPartition(false).build(); + given(partitionService.resolve(ServiceType.TB_CORE, tenantId, deviceId)).willReturn(notMyTpi); + + // WHEN + service.updateInactivityStateIfExpired(System.currentTimeMillis(), deviceId, stateData); + + // THEN + assertThat(service.deviceStates.get(deviceId)).isNull(); + assertThat(service.deviceStates.size()).isEqualTo(0); + assertThat(service.deviceStates.isEmpty()).isTrue(); + } + + @ParameterizedTest + @MethodSource("provideParametersForUpdateInactivityStateIfExpired") + public void givenTestParameters_whenUpdateInactivityStateIfExpired_thenShouldBeInTheExpectedStateAndPerformExpectedActions( + boolean activityState, long ts, long lastActivityTime, long lastInactivityAlarmTime, long inactivityTimeout, long deviceCreationTime, + boolean expectedActivityState, long expectedLastInactivityAlarmTime, boolean shouldUpdateActivityStateToInactive + ) { + // GIVEN + var state = DeviceState.builder() + .active(activityState) + .lastActivityTime(lastActivityTime) + .lastInactivityAlarmTime(lastInactivityAlarmTime) + .inactivityTimeout(inactivityTimeout) + .build(); + + var deviceStateData = DeviceStateData.builder() + .tenantId(tenantId) + .deviceId(deviceId) + .deviceCreationTime(deviceCreationTime) + .metaData(new TbMsgMetaData()) + .state(state) + .build(); + + if (shouldUpdateActivityStateToInactive) { + given(partitionService.resolve(ServiceType.TB_CORE, tenantId, deviceId)).willReturn(tpi); + } + + // WHEN + service.updateInactivityStateIfExpired(ts, deviceId, deviceStateData); + + // THEN + assertThat(state.isActive()).isEqualTo(expectedActivityState); + assertThat(state.getLastInactivityAlarmTime()).isEqualTo(expectedLastInactivityAlarmTime); + + if (shouldUpdateActivityStateToInactive) { + then(telemetrySubscriptionService).should().saveAttrAndNotify( + eq(TenantId.SYS_TENANT_ID), eq(deviceId), eq(SERVER_SCOPE), eq(ACTIVITY_STATE), eq(false), any() + ); + + var msgCaptor = ArgumentCaptor.forClass(TbMsg.class); + then(clusterService).should().pushMsgToRuleEngine(eq(tenantId), eq(deviceId), msgCaptor.capture(), any()); + var actualMsg = msgCaptor.getValue(); + assertThat(actualMsg.getType()).isEqualTo(TbMsgType.INACTIVITY_EVENT.name()); + assertThat(actualMsg.getOriginator()).isEqualTo(deviceId); + + var notificationCaptor = ArgumentCaptor.forClass(DeviceActivityTrigger.class); + then(notificationRuleProcessor).should().process(notificationCaptor.capture()); + var actualNotification = notificationCaptor.getValue(); + assertThat(actualNotification.getTenantId()).isEqualTo(tenantId); + assertThat(actualNotification.getDeviceId()).isEqualTo(deviceId); + assertThat(actualNotification.isActive()).isFalse(); + + then(telemetrySubscriptionService).should().saveAttrAndNotify( + eq(TenantId.SYS_TENANT_ID), eq(deviceId), eq(SERVER_SCOPE), + eq(INACTIVITY_ALARM_TIME), eq(expectedLastInactivityAlarmTime), any() + ); + } + } + + private static Stream provideParametersForUpdateInactivityStateIfExpired() { + return Stream.of( + Arguments.of(false, 100, 70, 90, 70, 60, false, 90, false), + + Arguments.of(false, 100, 40, 50, 70, 10, false, 50, false), + + Arguments.of(false, 100, 25, 60, 75, 25, false, 60, false), + + Arguments.of(false, 100, 60, 70, 10, 50, false, 70, false), + + Arguments.of(false, 100, 10, 15, 90, 10, false, 15, false), + + Arguments.of(false, 100, 0, 40, 75, 0, false, 40, false), + + Arguments.of(true, 100, 90, 80, 80, 50, true, 80, false), + + Arguments.of(true, 100, 95, 90, 10, 50, true, 90, false), + + Arguments.of(true, 100, 10, 10, 90, 10, false, 100, true), + + Arguments.of(true, 100, 10, 10, 90, 11, true, 10, false), + + Arguments.of(true, 100, 15, 10, 85, 5, false, 100, true), + + Arguments.of(true, 100, 15, 10, 75, 5, false, 100, true), + + Arguments.of(true, 100, 95, 90, 5, 50, false, 100, true), + + Arguments.of(true, 100, 0, 0, 101, 0, true, 0, false), + + Arguments.of(true, 100, 0, 0, 100, 0, false, 100, true), + + Arguments.of(true, 100, 0, 0, 99, 0, false, 100, true), + + Arguments.of(true, 100, 0, 0, 120, 10, true, 0, false), + + Arguments.of(true, 100, 50, 0, 100, 0, true, 0, false), + + Arguments.of(true, 100, 10, 0, 91, 0, true, 0, false), + + Arguments.of(true, 100, 90, 0, 10, 0, false, 100, true), + + Arguments.of(true, 100, 100, 100, 1, 0, true, 100, false), + + Arguments.of(true, 100, 100, 100, 100, 100, true, 100, false), + + Arguments.of(false, 100, 59, 60, 30, 10, false, 60, false), + + Arguments.of(true, 100, 60, 60, 30, 10, false, 100, true), + + Arguments.of(true, 100, 61, 60, 30, 10, false, 100, true), + + Arguments.of(true, 0, 0, 0, 1, 0, true, 0, false), + + Arguments.of(true, 0, 0, 0, 0, 0, false, 0, true), + + Arguments.of(true, 100, 90, 80, 20, 70, true, 80, false), + + Arguments.of(true, 100, 80, 90, 30, 70, true, 90, false) + ); + } + +} From 5e6dd652d7fb85b45251d2d70427f8108a81b5ce Mon Sep 17 00:00:00 2001 From: Dmytro Skarzhynets Date: Thu, 7 Sep 2023 12:41:20 +0300 Subject: [PATCH 5/9] Adjust timeout passed since device creation time check to be not strict --- .../server/service/state/DefaultDeviceStateService.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/application/src/main/java/org/thingsboard/server/service/state/DefaultDeviceStateService.java b/application/src/main/java/org/thingsboard/server/service/state/DefaultDeviceStateService.java index f9afb5b2e5..23276dc104 100644 --- a/application/src/main/java/org/thingsboard/server/service/state/DefaultDeviceStateService.java +++ b/application/src/main/java/org/thingsboard/server/service/state/DefaultDeviceStateService.java @@ -487,7 +487,7 @@ public class DefaultDeviceStateService extends AbstractPartitionBasedService Date: Thu, 7 Sep 2023 15:58:19 +0300 Subject: [PATCH 6/9] Add decrease inactivity timeout test with a controlled time --- .../state/DefaultDeviceStateService.java | 20 +++++--- .../state/DefaultDeviceStateServiceTest.java | 51 +++++++++++++++++++ 2 files changed, 63 insertions(+), 8 deletions(-) diff --git a/application/src/main/java/org/thingsboard/server/service/state/DefaultDeviceStateService.java b/application/src/main/java/org/thingsboard/server/service/state/DefaultDeviceStateService.java index 23276dc104..1f9239cea3 100644 --- a/application/src/main/java/org/thingsboard/server/service/state/DefaultDeviceStateService.java +++ b/application/src/main/java/org/thingsboard/server/service/state/DefaultDeviceStateService.java @@ -223,7 +223,7 @@ public class DefaultDeviceStateService extends AbstractPartitionBasedService= deviceState.getLastActivityTime()) { deviceState.setLastInactivityAlarmTime(0L); @@ -425,7 +425,7 @@ public class DefaultDeviceStateService extends AbstractPartitionBasedService { log.debug("Calculating state updates. tpi {} for {} devices", tpi.getFullTopicName(), deviceIds.size()); Set idsFromRemovedTenant = new HashSet<>(); @@ -455,7 +455,7 @@ public class DefaultDeviceStateService extends AbstractPartitionBasedService> stats = new HashMap<>(); for (DeviceStateData stateData : deviceStates.values()) { Pair tenantDevicesActivity = stats.computeIfAbsent(stateData.getTenantId(), @@ -798,7 +798,7 @@ public class DefaultDeviceStateService extends AbstractPartitionBasedService(deviceId, key, value)); } else { tsSubService.saveAttrAndNotify(TenantId.SYS_TENANT_ID, deviceId, SERVER_SCOPE, key, value, new TelemetrySaveCallback<>(deviceId, key, value)); @@ -809,13 +809,17 @@ public class DefaultDeviceStateService extends AbstractPartitionBasedService(deviceId, key, value)); } else { tsSubService.saveAttrAndNotify(TenantId.SYS_TENANT_ID, deviceId, SERVER_SCOPE, key, value, new TelemetrySaveCallback<>(deviceId, key, value)); } } + long getCurrentTimeMillis() { + return System.currentTimeMillis(); + } + private static class TelemetrySaveCallback implements FutureCallback { private final DeviceId deviceId; private final String key; diff --git a/application/src/test/java/org/thingsboard/server/service/state/DefaultDeviceStateServiceTest.java b/application/src/test/java/org/thingsboard/server/service/state/DefaultDeviceStateServiceTest.java index b06bbab0de..d03771cde3 100644 --- a/application/src/test/java/org/thingsboard/server/service/state/DefaultDeviceStateServiceTest.java +++ b/application/src/test/java/org/thingsboard/server/service/state/DefaultDeviceStateServiceTest.java @@ -54,6 +54,7 @@ import java.util.Collections; import java.util.List; import java.util.Map; import java.util.UUID; +import java.util.concurrent.atomic.AtomicLong; import java.util.stream.Stream; import static org.assertj.core.api.Assertions.assertThat; @@ -356,6 +357,56 @@ public class DefaultDeviceStateServiceTest { verify(telemetrySubscriptionService, times(1)).saveAttrAndNotify(any(), eq(deviceId), any(), eq(ACTIVITY_STATE), eq(isActive), any()); } + @ParameterizedTest + @MethodSource("provideParametersForDecreaseInactivityTimeout") + public void givenTestParameters_whenOnDeviceInactivityTimeout_thenShouldBeInTheExpectedStateAndPerformExpectedActions( + boolean activityState, long newInactivityTimeout, long timeIncrement, boolean expectedActivityState + ) throws Exception { + // GIVEN + long defaultInactivityTimeout = 10000; + initStateService(defaultInactivityTimeout); + + var currentTime = new AtomicLong(System.currentTimeMillis()); + + DeviceState deviceState = DeviceState.builder() + .active(activityState) + .lastActivityTime(currentTime.get()) + .inactivityTimeout(defaultInactivityTimeout) + .build(); + + DeviceStateData deviceStateData = DeviceStateData.builder() + .tenantId(tenantId) + .deviceId(deviceId) + .state(deviceState) + .metaData(new TbMsgMetaData()) + .build(); + + service.deviceStates.put(deviceId, deviceStateData); + service.getPartitionedEntities(tpi).add(deviceId); + + given(service.getCurrentTimeMillis()).willReturn(currentTime.addAndGet(timeIncrement)); + + // WHEN + service.onDeviceInactivityTimeoutUpdate(tenantId, deviceId, newInactivityTimeout); + + // THEN + assertThat(deviceState.getInactivityTimeout()).isEqualTo(newInactivityTimeout); + assertThat(deviceState.isActive()).isEqualTo(expectedActivityState); + if (activityState && !expectedActivityState) { + then(telemetrySubscriptionService).should().saveAttrAndNotify( + any(), eq(deviceId), any(), eq(ACTIVITY_STATE), eq(false), any() + ); + } + } + + private static Stream provideParametersForDecreaseInactivityTimeout() { + return Stream.of( + Arguments.of(true, 1, 0, true), + + Arguments.of(true, 1, 1, false) + ); + } + @Test public void givenStateDataIsNull_whenUpdateInactivityTimeoutIfExpired_thenShouldCleanupDevice() { // GIVEN From a4d4e8c6eb6f621ab441454b2b40c9fc026c348a Mon Sep 17 00:00:00 2001 From: Dmytro Skarzhynets Date: Thu, 7 Sep 2023 15:59:03 +0300 Subject: [PATCH 7/9] Remove redundant sleep --- .../server/service/state/DefaultDeviceStateServiceTest.java | 1 - 1 file changed, 1 deletion(-) diff --git a/application/src/test/java/org/thingsboard/server/service/state/DefaultDeviceStateServiceTest.java b/application/src/test/java/org/thingsboard/server/service/state/DefaultDeviceStateServiceTest.java index d03771cde3..b9716320a6 100644 --- a/application/src/test/java/org/thingsboard/server/service/state/DefaultDeviceStateServiceTest.java +++ b/application/src/test/java/org/thingsboard/server/service/state/DefaultDeviceStateServiceTest.java @@ -347,7 +347,6 @@ public class DefaultDeviceStateServiceTest { reset(telemetrySubscriptionService); long newTimeout = 1; - Thread.sleep(newTimeout); service.onDeviceInactivityTimeoutUpdate(tenantId, deviceId, newTimeout); verify(telemetrySubscriptionService, never()).saveAttrAndNotify(any(), eq(deviceId), any(), eq(ACTIVITY_STATE), any(), any()); From babc8844e9a552b910ace2779327a1cbef4ad623 Mon Sep 17 00:00:00 2001 From: Dmytro Skarzhynets Date: Thu, 14 Sep 2023 12:58:04 +0300 Subject: [PATCH 8/9] Add last reported activity in the past fix, cover with a test --- .../state/DefaultDeviceStateService.java | 4 + .../state/DefaultDeviceStateServiceTest.java | 102 ++++++++++++++++++ 2 files changed, 106 insertions(+) diff --git a/application/src/main/java/org/thingsboard/server/service/state/DefaultDeviceStateService.java b/application/src/main/java/org/thingsboard/server/service/state/DefaultDeviceStateService.java index 1f9239cea3..41a41f5e82 100644 --- a/application/src/main/java/org/thingsboard/server/service/state/DefaultDeviceStateService.java +++ b/application/src/main/java/org/thingsboard/server/service/state/DefaultDeviceStateService.java @@ -250,6 +250,10 @@ public class DefaultDeviceStateService extends AbstractPartitionBasedService provideParametersForUpdateActivityState() { + return Stream.of( + Arguments.of(true, 100, 120, 80, 80, false, false), + + Arguments.of(true, 100, 120, 100, 100, false, false), + + Arguments.of(false, 100, 120, 110, 110, false, true), + + + Arguments.of(true, 100, 100, 80, 80, false, false), + + Arguments.of(true, 100, 100, 100, 100, false, false), + + Arguments.of(false, 100, 100, 110, 0, true, true), + + + Arguments.of(false, 100, 110, 110, 0, true, true), + + Arguments.of(false, 100, 110, 120, 0, true, true), + + + Arguments.of(true, 0, 0, 0, 0, false, false), + + Arguments.of(false, 0, 0, 0, 0, true, true) + ); + } + @ParameterizedTest @MethodSource("provideParametersForDecreaseInactivityTimeout") public void givenTestParameters_whenOnDeviceInactivityTimeout_thenShouldBeInTheExpectedStateAndPerformExpectedActions( From 65a4f002bbbdcc9f1a057a43f9fdd82902923ae0 Mon Sep 17 00:00:00 2001 From: Dmytro Skarzhynets Date: Thu, 14 Sep 2023 15:49:48 +0300 Subject: [PATCH 9/9] Fix copy-paste error --- .../server/service/state/DefaultDeviceStateServiceTest.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/application/src/test/java/org/thingsboard/server/service/state/DefaultDeviceStateServiceTest.java b/application/src/test/java/org/thingsboard/server/service/state/DefaultDeviceStateServiceTest.java index 06cde65c36..746da9a8f0 100644 --- a/application/src/test/java/org/thingsboard/server/service/state/DefaultDeviceStateServiceTest.java +++ b/application/src/test/java/org/thingsboard/server/service/state/DefaultDeviceStateServiceTest.java @@ -363,7 +363,7 @@ public class DefaultDeviceStateServiceTest { service.deviceStates.put(deviceId, deviceStateDataMock); // WHEN - service.updateActivityState(deviceId, deviceStateDataMock, System.currentTimeMillis()); + service.updateActivityState(deviceId, null, System.currentTimeMillis()); // THEN assertThat(service.deviceStates.get(deviceId)).isNull();