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..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 @@ -223,7 +223,7 @@ public class DefaultDeviceStateService extends AbstractPartitionBasedService= deviceState.getLastActivityTime()) { deviceState.setLastInactivityAlarmTime(0L); @@ -425,7 +429,7 @@ public class DefaultDeviceStateService extends AbstractPartitionBasedService { log.debug("Calculating state updates. tpi {} for {} devices", tpi.getFullTopicName(), deviceIds.size()); Set idsFromRemovedTenant = new HashSet<>(); @@ -455,7 +459,7 @@ public class DefaultDeviceStateService extends AbstractPartitionBasedService> stats = new HashMap<>(); for (DeviceStateData stateData : deviceStates.values()) { Pair tenantDevicesActivity = stats.computeIfAbsent(stateData.getTenantId(), @@ -486,8 +490,8 @@ 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 +813,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 adb46affe8..05624b8741 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,22 +15,27 @@ */ 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.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.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; 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; @@ -49,17 +54,28 @@ import java.util.Collections; 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 java.util.concurrent.atomic.AtomicLong; +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.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.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; +import static org.thingsboard.server.service.state.DefaultDeviceStateService.LAST_ACTIVITY_TIME; -@RunWith(MockitoJUnitRunner.class) +@ExtendWith(MockitoExtension.class) public class DefaultDeviceStateServiceTest { @Mock @@ -76,6 +92,12 @@ public class DefaultDeviceStateServiceTest { DeviceStateData deviceStateDataMock; @Mock EntityQueryRepository entityQueryRepository; + @Mock + TelemetrySubscriptionService telemetrySubscriptionService; + @Mock + NotificationRuleProcessor notificationRuleProcessor; + @Mock + DefaultTbApiUsageReportClient defaultTbApiUsageReportClient; TenantId tenantId = new TenantId(UUID.fromString("00797a3b-7aeb-4b5b-b57a-c2a810d0f112")); DeviceId deviceId = DeviceId.fromString("00797a3b-7aeb-4b5b-b57a-c2a810d0f112"); @@ -83,31 +105,23 @@ public class DefaultDeviceStateServiceTest { DefaultDeviceStateService service; - TelemetrySubscriptionService telemetrySubscriptionService; - - @Before + @BeforeEach public void setUp() { - service = spy(new DefaultDeviceStateService(deviceService, attributesService, tsService, clusterService, partitionService, entityQueryRepository, null, mock(DefaultTbApiUsageReportClient.class), mock(NotificationRuleProcessor.class))); - telemetrySubscriptionService = Mockito.mock(TelemetrySubscriptionService.class); + service = spy(new DefaultDeviceStateService(deviceService, attributesService, tsService, clusterService, partitionService, entityQueryRepository, null, defaultTbApiUsageReportClient, 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 @@ -115,8 +129,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 @@ -152,14 +166,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, ServiceType.TB_CORE, Map.of( new QueueKey(ServiceType.TB_CORE), Collections.singleton(tpi) )); @@ -188,7 +207,7 @@ public class DefaultDeviceStateServiceTest { service.checkStates(); activityVerify(false); - Mockito.reset(telemetrySubscriptionService); + reset(telemetrySubscriptionService); long increase = 100; long newTimeout = System.currentTimeMillis() - deviceState.getLastActivityTime() + increase; @@ -199,7 +218,7 @@ public class DefaultDeviceStateServiceTest { service.checkStates(); activityVerify(false); - Mockito.reset(telemetrySubscriptionService); + reset(telemetrySubscriptionService); service.onDeviceActivity(tenantId, deviceId, System.currentTimeMillis()); activityVerify(true); @@ -226,18 +245,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); @@ -267,11 +286,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 @@ -292,16 +311,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); @@ -330,17 +347,336 @@ 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()); + } + + @Test + public void givenStateDataIsNull_whenUpdateActivityState_thenShouldCleanupDevice() { + // GIVEN + service.deviceStates.put(deviceId, deviceStateDataMock); + + // WHEN + service.updateActivityState(deviceId, null, System.currentTimeMillis()); + + // THEN + assertThat(service.deviceStates.get(deviceId)).isNull(); + assertThat(service.deviceStates.size()).isEqualTo(0); + assertThat(service.deviceStates.isEmpty()).isTrue(); + } + + + @ParameterizedTest + @MethodSource("provideParametersForUpdateActivityState") + public void givenTestParameters_whenUpdateActivityState_thenShouldBeInTheExpectedStateAndPerformExpectedActions( + boolean activityState, long previousActivityTime, long lastReportedActivity, long inactivityAlarmTime, + long expectedInactivityAlarmTime, boolean shouldSetInactivityAlarmTimeToZero, + boolean shouldUpdateActivityStateToActive + ) { + // GIVEN + DeviceState deviceState = DeviceState.builder() + .active(activityState) + .lastActivityTime(previousActivityTime) + .lastInactivityAlarmTime(inactivityAlarmTime) + .inactivityTimeout(10000) + .build(); + + DeviceStateData deviceStateData = DeviceStateData.builder() + .tenantId(tenantId) + .deviceId(deviceId) + .state(deviceState) + .metaData(new TbMsgMetaData()) + .build(); + + // WHEN + service.updateActivityState(deviceId, deviceStateData, lastReportedActivity); + + // THEN + assertThat(deviceState.isActive()).isEqualTo(true); + assertThat(deviceState.getLastActivityTime()).isEqualTo(lastReportedActivity); + then(telemetrySubscriptionService).should().saveAttrAndNotify( + any(), eq(deviceId), any(), eq(LAST_ACTIVITY_TIME), eq(lastReportedActivity), any() + ); + + assertThat(deviceState.getLastInactivityAlarmTime()).isEqualTo(expectedInactivityAlarmTime); + if (shouldSetInactivityAlarmTimeToZero) { + then(telemetrySubscriptionService).should().saveAttrAndNotify( + any(), eq(deviceId), any(), eq(INACTIVITY_ALARM_TIME), eq(0L), any() + ); + } + + if (shouldUpdateActivityStateToActive) { + then(telemetrySubscriptionService).should().saveAttrAndNotify( + eq(TenantId.SYS_TENANT_ID), eq(deviceId), eq(SERVER_SCOPE), eq(ACTIVITY_STATE), eq(true), 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.ACTIVITY_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()).isTrue(); + } + } + + private static Stream 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( + 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 + 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) + ); } -} \ No newline at end of file +}