From 8dc9a68c625b34ad7523aea6bdfe5817fffc4b12 Mon Sep 17 00:00:00 2001 From: Dmytro Skarzhynets Date: Fri, 6 Jun 2025 18:47:45 +0300 Subject: [PATCH] Update cached activity status only after a successful database save --- .../state/DefaultDeviceStateService.java | 97 ++-- .../DefaultTelemetrySubscriptionService.java | 4 +- .../telemetry/InternalTelemetryService.java | 4 +- .../src/main/resources/thingsboard.yml | 2 + .../state/DefaultDeviceStateServiceTest.java | 527 ++++++++---------- 5 files changed, 283 insertions(+), 351 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 cc476d377d..f46323f702 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 @@ -27,11 +27,10 @@ import jakarta.annotation.Nonnull; import jakarta.annotation.Nullable; import jakarta.annotation.PostConstruct; import jakarta.annotation.PreDestroy; -import lombok.Getter; import lombok.RequiredArgsConstructor; -import lombok.Setter; import lombok.extern.slf4j.Slf4j; import org.apache.commons.lang3.tuple.Pair; +import org.checkerframework.checker.nullness.qual.NonNull; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.beans.factory.annotation.Value; import org.springframework.context.annotation.Lazy; @@ -170,35 +169,22 @@ public class DefaultDeviceStateService extends AbstractPartitionBasedService> stats = new HashMap<>(); for (DeviceStateData stateData : deviceStates.values()) { @@ -587,13 +572,12 @@ public class DefaultDeviceStateService extends AbstractPartitionBasedService { + stateData.getState().setActive(active); + pushRuleEngineMessage(stateData, active ? TbMsgType.ACTIVITY_EVENT : TbMsgType.INACTIVITY_EVENT); + TbMsgMetaData metaData = stateData.getMetaData(); + notificationRuleProcessor.process(DeviceActivityTrigger.builder() + .tenantId(tenantId) + .customerId(stateData.getCustomerId()) + .deviceId(deviceId) + .active(active) + .deviceName(metaData.getValue("deviceName")) + .deviceType(metaData.getValue("deviceType")) + .deviceLabel(metaData.getValue("deviceLabel")) + .build()); + }, deviceStateCallbackExecutor); } boolean cleanDeviceStateIfBelongsToExternalPartition(TenantId tenantId, final DeviceId deviceId) { @@ -634,8 +625,7 @@ public class DefaultDeviceStateService extends AbstractPartitionBasedService save(TenantId tenantId, DeviceId deviceId, String key, long value) { + return save(tenantId, deviceId, new LongDataEntry(key, value), getCurrentTimeMillis()); } - private void save(TenantId tenantId, DeviceId deviceId, String key, boolean value) { - save(tenantId, deviceId, new BooleanDataEntry(key, value), getCurrentTimeMillis()); + private ListenableFuture save(TenantId tenantId, DeviceId deviceId, String key, boolean value) { + return save(tenantId, deviceId, new BooleanDataEntry(key, value), getCurrentTimeMillis()); } - private void save(TenantId tenantId, DeviceId deviceId, KvEntry kvEntry, long ts) { + private ListenableFuture save(TenantId tenantId, DeviceId deviceId, KvEntry kvEntry, long ts) { + ListenableFuture future; if (persistToTelemetry) { - tsSubService.saveTimeseriesInternal(TimeseriesSaveRequest.builder() + future = tsSubService.saveTimeseriesInternal(TimeseriesSaveRequest.builder() .tenantId(tenantId) .entityId(deviceId) .entry(new BasicTsKvEntry(ts, kvEntry)) @@ -895,7 +889,7 @@ public class DefaultDeviceStateService extends AbstractPartitionBasedService(deviceId, kvEntry)) .build()); } else { - tsSubService.saveAttributes(AttributesSaveRequest.builder() + future = tsSubService.saveAttributesInternal(AttributesSaveRequest.builder() .tenantId(tenantId) .entityId(deviceId) .scope(AttributeScope.SERVER_SCOPE) @@ -903,20 +897,14 @@ public class DefaultDeviceStateService extends AbstractPartitionBasedService(deviceId, kvEntry)) .build()); } + return Futures.transform(future, __ -> null, MoreExecutors.directExecutor()); } long getCurrentTimeMillis() { return System.currentTimeMillis(); } - private static class TelemetrySaveCallback implements FutureCallback { - private final DeviceId deviceId; - private final KvEntry kvEntry; - - TelemetrySaveCallback(DeviceId deviceId, KvEntry kvEntry) { - this.deviceId = deviceId; - this.kvEntry = kvEntry; - } + private record TelemetrySaveCallback(DeviceId deviceId, KvEntry kvEntry) implements FutureCallback { @Override public void onSuccess(@Nullable T result) { @@ -924,9 +912,10 @@ public class DefaultDeviceStateService extends AbstractPartitionBasedService> saveAttributesInternal(AttributesSaveRequest request) { TenantId tenantId = request.getTenantId(); EntityId entityId = request.getEntityId(); AttributesSaveRequest.Strategy strategy = request.getStrategy(); @@ -228,6 +227,7 @@ public class DefaultTelemetrySubscriptionService extends AbstractSubscriptionSer if (strategy.sendWsUpdate()) { addWsCallback(resultFuture, success -> onAttributesUpdate(tenantId, entityId, request.getScope().name(), request.getEntries())); } + return resultFuture; } private static boolean shouldSendSharedAttributesUpdatedNotification(AttributesSaveRequest request) { diff --git a/application/src/main/java/org/thingsboard/server/service/telemetry/InternalTelemetryService.java b/application/src/main/java/org/thingsboard/server/service/telemetry/InternalTelemetryService.java index 8a76aa1d14..79f0beab41 100644 --- a/application/src/main/java/org/thingsboard/server/service/telemetry/InternalTelemetryService.java +++ b/application/src/main/java/org/thingsboard/server/service/telemetry/InternalTelemetryService.java @@ -23,6 +23,8 @@ import org.thingsboard.rule.engine.api.TimeseriesDeleteRequest; import org.thingsboard.rule.engine.api.TimeseriesSaveRequest; import org.thingsboard.server.common.data.kv.TimeseriesSaveResult; +import java.util.List; + /** * Created by ashvayka on 27.03.18. */ @@ -30,7 +32,7 @@ public interface InternalTelemetryService extends RuleEngineTelemetryService { ListenableFuture saveTimeseriesInternal(TimeseriesSaveRequest request); - void saveAttributesInternal(AttributesSaveRequest request); + ListenableFuture> saveAttributesInternal(AttributesSaveRequest request); void deleteTimeseriesInternal(TimeseriesDeleteRequest request); diff --git a/application/src/main/resources/thingsboard.yml b/application/src/main/resources/thingsboard.yml index 6d27b9bd7e..09820a81cd 100644 --- a/application/src/main/resources/thingsboard.yml +++ b/application/src/main/resources/thingsboard.yml @@ -898,6 +898,8 @@ state: # Used only when state.persistToTelemetry is set to 'true' and Cassandra is used for timeseries data. # 0 means time-to-live mechanism is disabled. telemetryTtl: "${STATE_TELEMETRY_TTL:0}" + # Number of device records to fetch per batch when initializing device activity states + initFetchPackSize: "${TB_DEVICE_STATE_INIT_FETCH_PACK_SIZE:50000}" # Configuration properties for rule nodes related to device activity state rule: node: 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 26e913eacc..cbf7363441 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 @@ -16,6 +16,9 @@ package org.thingsboard.server.service.state; import com.google.common.util.concurrent.Futures; +import com.google.common.util.concurrent.ListeningExecutorService; +import com.google.common.util.concurrent.MoreExecutors; +import org.junit.jupiter.api.AfterEach; import org.junit.jupiter.api.BeforeEach; import org.junit.jupiter.api.Test; import org.junit.jupiter.api.extension.ExtendWith; @@ -31,13 +34,12 @@ import org.thingsboard.rule.engine.api.AttributesSaveRequest; import org.thingsboard.server.cluster.TbClusterService; import org.thingsboard.server.common.data.AttributeScope; import org.thingsboard.server.common.data.Device; -import org.thingsboard.server.common.data.DeviceIdInfo; import org.thingsboard.server.common.data.id.DeviceId; import org.thingsboard.server.common.data.id.EntityId; import org.thingsboard.server.common.data.id.TenantId; +import org.thingsboard.server.common.data.kv.AttributeKvEntry; 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.msg.TbMsg; import org.thingsboard.server.common.msg.TbMsgMetaData; import org.thingsboard.server.common.msg.notification.NotificationRuleProcessor; @@ -50,20 +52,22 @@ import org.thingsboard.server.dao.sql.query.EntityQueryRepository; import org.thingsboard.server.dao.timeseries.TimeseriesService; import org.thingsboard.server.gen.transport.TransportProtos; import org.thingsboard.server.queue.discovery.PartitionService; -import org.thingsboard.server.queue.discovery.QueueKey; -import org.thingsboard.server.queue.discovery.event.PartitionChangeEvent; import org.thingsboard.server.queue.usagestats.DefaultTbApiUsageReportClient; import org.thingsboard.server.service.telemetry.TelemetrySubscriptionService; +import java.time.Duration; import java.util.Collections; +import java.util.HashSet; import java.util.List; -import java.util.Map; +import java.util.Set; import java.util.UUID; +import java.util.concurrent.ConcurrentHashMap; +import java.util.concurrent.ConcurrentMap; import java.util.concurrent.CountDownLatch; import java.util.concurrent.ExecutorService; import java.util.concurrent.Executors; +import java.util.concurrent.ThreadLocalRandom; import java.util.concurrent.TimeUnit; -import java.util.concurrent.atomic.AtomicLong; import java.util.stream.Stream; import static org.assertj.core.api.Assertions.assertThat; @@ -77,8 +81,8 @@ import static org.mockito.BDDMockito.then; import static org.mockito.BDDMockito.willReturn; import static org.mockito.Mockito.doAnswer; import static org.mockito.Mockito.doReturn; +import static org.mockito.Mockito.lenient; 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; @@ -90,7 +94,10 @@ import static org.thingsboard.server.service.state.DefaultDeviceStateService.LAS import static org.thingsboard.server.service.state.DefaultDeviceStateService.LAST_DISCONNECT_TIME; @ExtendWith(MockitoExtension.class) -public class DefaultDeviceStateServiceTest { +class DefaultDeviceStateServiceTest { + + ListeningExecutorService deviceStateExecutor; + ListeningExecutorService deviceStateCallbackExecutor; @Mock DeviceService deviceService; @@ -113,25 +120,48 @@ public class DefaultDeviceStateServiceTest { @Mock DefaultTbApiUsageReportClient defaultTbApiUsageReportClient; - TenantId tenantId = new TenantId(UUID.fromString("00797a3b-7aeb-4b5b-b57a-c2a810d0f112")); - DeviceId deviceId = DeviceId.fromString("00797a3b-7aeb-4b5b-b57a-c2a810d0f112"); - TopicPartitionInfo tpi; + long defaultInactivityTimeoutMs = Duration.ofMinutes(10L).toMillis(); + + TenantId tenantId = TenantId.fromUUID(UUID.fromString("00797a3b-7aeb-4b5b-b57a-c2a810d0f112")); + DeviceId deviceId = DeviceId.fromString("c209f718-42e5-11f0-9fe2-0242ac120002"); + TopicPartitionInfo tpi = TopicPartitionInfo.builder() + .topic("tb_core") + .partition(0) + .myPartition(true) + .build(); DefaultDeviceStateService service; @BeforeEach - public void setUp() { + void setUp() { service = spy(new DefaultDeviceStateService(deviceService, attributesService, tsService, clusterService, partitionService, entityQueryRepository, null, defaultTbApiUsageReportClient, notificationRuleProcessor)); ReflectionTestUtils.setField(service, "tsSubService", telemetrySubscriptionService); + ReflectionTestUtils.setField(service, "defaultInactivityTimeoutMs", defaultInactivityTimeoutMs); ReflectionTestUtils.setField(service, "defaultStateCheckIntervalInSec", 60); ReflectionTestUtils.setField(service, "defaultActivityStatsIntervalInSec", 60); - ReflectionTestUtils.setField(service, "initFetchPackSize", 10); + ReflectionTestUtils.setField(service, "initFetchPackSize", 50000); + + deviceStateExecutor = MoreExecutors.newDirectExecutorService(); + ReflectionTestUtils.setField(service, "deviceStateExecutor", deviceStateExecutor); + + deviceStateCallbackExecutor = MoreExecutors.newDirectExecutorService(); + ReflectionTestUtils.setField(service, "deviceStateCallbackExecutor", deviceStateCallbackExecutor); + + lenient().when(partitionService.resolve(ServiceType.TB_CORE, tenantId, deviceId)).thenReturn(tpi); + + ConcurrentMap> partitionedEntities = new ConcurrentHashMap<>(); + partitionedEntities.put(tpi, new HashSet<>()); + ReflectionTestUtils.setField(service, "partitionedEntities", partitionedEntities); + } - tpi = TopicPartitionInfo.builder().myPartition(true).build(); + @AfterEach + void cleanup() { + deviceStateExecutor.shutdownNow(); + deviceStateCallbackExecutor.shutdownNow(); } @Test - public void givenDeviceBelongsToExternalPartition_whenOnDeviceConnect_thenCleansStateAndDoesNotReportConnect() { + void givenDeviceBelongsToExternalPartition_whenOnDeviceConnect_thenCleansStateAndDoesNotReportConnect() { // GIVEN doReturn(true).when(service).cleanDeviceStateIfBelongsToExternalPartition(tenantId, deviceId); @@ -149,7 +179,7 @@ public class DefaultDeviceStateServiceTest { @ParameterizedTest @ValueSource(longs = {Long.MIN_VALUE, -100, -1}) - public void givenNegativeLastConnectTime_whenOnDeviceConnect_thenSkipsThisEvent(long negativeLastConnectTime) { + void givenNegativeLastConnectTime_whenOnDeviceConnect_thenSkipsThisEvent(long negativeLastConnectTime) { // GIVEN doReturn(false).when(service).cleanDeviceStateIfBelongsToExternalPartition(tenantId, deviceId); @@ -166,7 +196,7 @@ public class DefaultDeviceStateServiceTest { @ParameterizedTest @MethodSource("provideOutdatedTimestamps") - public void givenOutdatedLastConnectTime_whenOnDeviceDisconnect_thenSkipsThisEvent(long outdatedLastConnectTime, long currentLastConnectTime) { + void givenOutdatedLastConnectTime_whenOnDeviceDisconnect_thenSkipsThisEvent(long outdatedLastConnectTime, long currentLastConnectTime) { // GIVEN doReturn(false).when(service).cleanDeviceStateIfBelongsToExternalPartition(tenantId, deviceId); @@ -188,7 +218,7 @@ public class DefaultDeviceStateServiceTest { } @Test - public void givenDeviceBelongsToMyPartition_whenOnDeviceConnect_thenReportsConnect() { + void givenDeviceBelongsToMyPartition_whenOnDeviceConnect_thenReportsConnect() { // GIVEN var deviceStateData = DeviceStateData.builder() .tenantId(tenantId) @@ -202,11 +232,13 @@ public class DefaultDeviceStateServiceTest { service.deviceStates.put(deviceId, deviceStateData); long lastConnectTime = System.currentTimeMillis(); + mockSuccessfulSaveAttributes(); + // WHEN service.onDeviceConnect(tenantId, deviceId, lastConnectTime); // THEN - then(telemetrySubscriptionService).should().saveAttributes(argThat(request -> + then(telemetrySubscriptionService).should().saveAttributesInternal(argThat(request -> request.getTenantId().equals(tenantId) && request.getEntityId().equals(deviceId) && request.getScope().equals(AttributeScope.SERVER_SCOPE) && request.getEntries().get(0).getKey().equals(LAST_CONNECT_TIME) && @@ -221,7 +253,7 @@ public class DefaultDeviceStateServiceTest { } @Test - public void givenDeviceBelongsToExternalPartition_whenOnDeviceDisconnect_thenCleansStateAndDoesNotReportDisconnect() { + void givenDeviceBelongsToExternalPartition_whenOnDeviceDisconnect_thenCleansStateAndDoesNotReportDisconnect() { // GIVEN doReturn(true).when(service).cleanDeviceStateIfBelongsToExternalPartition(tenantId, deviceId); @@ -238,7 +270,7 @@ public class DefaultDeviceStateServiceTest { @ParameterizedTest @ValueSource(longs = {Long.MIN_VALUE, -100, -1}) - public void givenNegativeLastDisconnectTime_whenOnDeviceDisconnect_thenSkipsThisEvent(long negativeLastDisconnectTime) { + void givenNegativeLastDisconnectTime_whenOnDeviceDisconnect_thenSkipsThisEvent(long negativeLastDisconnectTime) { // GIVEN doReturn(false).when(service).cleanDeviceStateIfBelongsToExternalPartition(tenantId, deviceId); @@ -254,7 +286,7 @@ public class DefaultDeviceStateServiceTest { @ParameterizedTest @MethodSource("provideOutdatedTimestamps") - public void givenOutdatedLastDisconnectTime_whenOnDeviceDisconnect_thenSkipsThisEvent(long outdatedLastDisconnectTime, long currentLastDisconnectTime) { + void givenOutdatedLastDisconnectTime_whenOnDeviceDisconnect_thenSkipsThisEvent(long outdatedLastDisconnectTime, long currentLastDisconnectTime) { // GIVEN doReturn(false).when(service).cleanDeviceStateIfBelongsToExternalPartition(tenantId, deviceId); @@ -275,7 +307,7 @@ public class DefaultDeviceStateServiceTest { } @Test - public void givenDeviceBelongsToMyPartition_whenOnDeviceDisconnect_thenReportsDisconnect() { + void givenDeviceBelongsToMyPartition_whenOnDeviceDisconnect_thenReportsDisconnect() { // GIVEN var deviceStateData = DeviceStateData.builder() .tenantId(tenantId) @@ -289,11 +321,13 @@ public class DefaultDeviceStateServiceTest { service.deviceStates.put(deviceId, deviceStateData); long lastDisconnectTime = System.currentTimeMillis(); + mockSuccessfulSaveAttributes(); + // WHEN service.onDeviceDisconnect(tenantId, deviceId, lastDisconnectTime); // THEN - then(telemetrySubscriptionService).should().saveAttributes(argThat(request -> + then(telemetrySubscriptionService).should().saveAttributesInternal(argThat(request -> request.getTenantId().equals(tenantId) && request.getEntityId().equals(deviceId) && request.getScope().equals(AttributeScope.SERVER_SCOPE) && request.getEntries().get(0).getKey().equals(LAST_DISCONNECT_TIME) && @@ -308,7 +342,7 @@ public class DefaultDeviceStateServiceTest { } @Test - public void givenDeviceBelongsToExternalPartition_whenOnDeviceInactivity_thenCleansStateAndDoesNotReportInactivity() { + void givenDeviceBelongsToExternalPartition_whenOnDeviceInactivity_thenCleansStateAndDoesNotReportInactivity() { // GIVEN doReturn(true).when(service).cleanDeviceStateIfBelongsToExternalPartition(tenantId, deviceId); @@ -325,7 +359,7 @@ public class DefaultDeviceStateServiceTest { @ParameterizedTest @ValueSource(longs = {Long.MIN_VALUE, -100, -1}) - public void givenNegativeLastInactivityTime_whenOnDeviceInactivity_thenSkipsThisEvent(long negativeLastInactivityTime) { + void givenNegativeLastInactivityTime_whenOnDeviceInactivity_thenSkipsThisEvent(long negativeLastInactivityTime) { // GIVEN doReturn(false).when(service).cleanDeviceStateIfBelongsToExternalPartition(tenantId, deviceId); @@ -341,7 +375,7 @@ public class DefaultDeviceStateServiceTest { @ParameterizedTest @MethodSource("provideOutdatedTimestamps") - public void givenReceivedInactivityTimeIsLessThanOrEqualToCurrentInactivityTime_whenOnDeviceInactivity_thenSkipsThisEvent( + void givenReceivedInactivityTimeIsLessThanOrEqualToCurrentInactivityTime_whenOnDeviceInactivity_thenSkipsThisEvent( long outdatedLastInactivityTime, long currentLastInactivityTime ) { // GIVEN @@ -365,7 +399,7 @@ public class DefaultDeviceStateServiceTest { @ParameterizedTest @MethodSource("provideOutdatedTimestamps") - public void givenReceivedInactivityTimeIsLessThanOrEqualToCurrentActivityTime_whenOnDeviceInactivity_thenSkipsThisEvent( + void givenReceivedInactivityTimeIsLessThanOrEqualToCurrentActivityTime_whenOnDeviceInactivity_thenSkipsThisEvent( long outdatedLastInactivityTime, long currentLastActivityTime ) { // GIVEN @@ -398,7 +432,7 @@ public class DefaultDeviceStateServiceTest { } @Test - public void givenDeviceBelongsToMyPartition_whenOnDeviceInactivity_thenReportsInactivity() { + void givenDeviceBelongsToMyPartition_whenOnDeviceInactivity_thenReportsInactivity() { // GIVEN var deviceStateData = DeviceStateData.builder() .tenantId(tenantId) @@ -412,17 +446,19 @@ public class DefaultDeviceStateServiceTest { service.deviceStates.put(deviceId, deviceStateData); long lastInactivityTime = System.currentTimeMillis(); + mockSuccessfulSaveAttributes(); + // WHEN service.onDeviceInactivity(tenantId, deviceId, lastInactivityTime); // THEN - then(telemetrySubscriptionService).should().saveAttributes(argThat(request -> + then(telemetrySubscriptionService).should().saveAttributesInternal(argThat(request -> request.getTenantId().equals(tenantId) && request.getEntityId().equals(deviceId) && request.getScope().equals(AttributeScope.SERVER_SCOPE) && request.getEntries().get(0).getKey().equals(INACTIVITY_ALARM_TIME) && request.getEntries().get(0).getValue().equals(lastInactivityTime) )); - then(telemetrySubscriptionService).should().saveAttributes(argThat(request -> + then(telemetrySubscriptionService).should().saveAttributesInternal(argThat(request -> request.getTenantId().equals(tenantId) && request.getEntityId().equals(deviceId) && request.getScope().equals(AttributeScope.SERVER_SCOPE) && request.getEntries().get(0).getKey().equals(ACTIVITY_STATE) && @@ -445,7 +481,7 @@ public class DefaultDeviceStateServiceTest { } @Test - public void givenInactivityTimeoutReached_whenUpdateInactivityStateIfExpired_thenReportsInactivity() { + void givenInactivityTimeoutReached_whenUpdateInactivityStateIfExpired_thenReportsInactivity() { // GIVEN var deviceStateData = DeviceStateData.builder() .tenantId(tenantId) @@ -456,16 +492,18 @@ public class DefaultDeviceStateServiceTest { given(partitionService.resolve(ServiceType.TB_CORE, tenantId, deviceId)).willReturn(tpi); + mockSuccessfulSaveAttributes(); + // WHEN service.updateInactivityStateIfExpired(System.currentTimeMillis(), deviceId, deviceStateData); // THEN - then(telemetrySubscriptionService).should().saveAttributes(argThat(request -> + then(telemetrySubscriptionService).should().saveAttributesInternal(argThat(request -> request.getTenantId().equals(tenantId) && request.getEntityId().equals(deviceId) && request.getScope().equals(AttributeScope.SERVER_SCOPE) && request.getEntries().get(0).getKey().equals(INACTIVITY_ALARM_TIME) )); - then(telemetrySubscriptionService).should().saveAttributes(argThat(request -> + then(telemetrySubscriptionService).should().saveAttributesInternal(argThat(request -> request.getTenantId().equals(tenantId) && request.getEntityId().equals(deviceId) && request.getScope().equals(AttributeScope.SERVER_SCOPE) && request.getEntries().get(0).getKey().equals(ACTIVITY_STATE) && @@ -488,7 +526,7 @@ public class DefaultDeviceStateServiceTest { } @Test - public void givenDeviceIdFromDeviceStatesMap_whenGetOrFetchDeviceStateData_thenNoStackOverflow() { + void givenDeviceIdFromDeviceStatesMap_whenGetOrFetchDeviceStateData_thenNoStackOverflow() { service.deviceStates.put(deviceId, deviceStateDataMock); DeviceStateData deviceStateData = service.getOrFetchDeviceStateData(deviceId); assertThat(deviceStateData).isEqualTo(deviceStateDataMock); @@ -496,7 +534,7 @@ public class DefaultDeviceStateServiceTest { } @Test - public void givenDeviceIdWithoutDeviceStateInMap_whenGetOrFetchDeviceStateData_thenFetchDeviceStateData() { + void givenDeviceIdWithoutDeviceStateInMap_whenGetOrFetchDeviceStateData_thenFetchDeviceStateData() { service.deviceStates.clear(); willReturn(deviceStateDataMock).given(service).fetchDeviceStateDataUsingSeparateRequests(deviceId); DeviceStateData deviceStateData = service.getOrFetchDeviceStateData(deviceId); @@ -504,172 +542,18 @@ public class DefaultDeviceStateServiceTest { verify(service).fetchDeviceStateDataUsingSeparateRequests(deviceId); } - private void initStateService(long timeout) throws InterruptedException { - service.stop(); - 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) - ), Collections.emptyMap()); - service.onApplicationEvent(event); - Thread.sleep(100); - } - - @Test - public void increaseInactivityForInactiveDeviceTest() throws Exception { - final long defaultTimeout = 1; - initStateService(defaultTimeout); - DeviceState deviceState = DeviceState.builder().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); - - service.onDeviceActivity(tenantId, deviceId, System.currentTimeMillis()); - activityVerify(true); - Thread.sleep(defaultTimeout); - service.checkStates(); - activityVerify(false); - - reset(telemetrySubscriptionService); - - long increase = 100; - long newTimeout = System.currentTimeMillis() - deviceState.getLastActivityTime() + increase; - - service.onDeviceInactivityTimeoutUpdate(tenantId, deviceId, newTimeout); - activityVerify(true); - Thread.sleep(increase); - service.checkStates(); - activityVerify(false); - - reset(telemetrySubscriptionService); - - service.onDeviceActivity(tenantId, deviceId, System.currentTimeMillis()); - activityVerify(true); - Thread.sleep(newTimeout + 5); - service.checkStates(); - activityVerify(false); - } - - @Test - public void increaseInactivityForActiveDeviceTest() throws Exception { - final long defaultTimeout = 1000; - initStateService(defaultTimeout); - DeviceState deviceState = DeviceState.builder().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); - - service.onDeviceActivity(tenantId, deviceId, System.currentTimeMillis()); - activityVerify(true); - - reset(telemetrySubscriptionService); - - long increase = 100; - long newTimeout = System.currentTimeMillis() - deviceState.getLastActivityTime() + increase; - - service.onDeviceInactivityTimeoutUpdate(tenantId, deviceId, newTimeout); - verify(telemetrySubscriptionService, never()).saveAttributes(argThat(request -> - request.getEntityId().equals(deviceId) && request.getEntries().get(0).getKey().equals(ACTIVITY_STATE) - )); - Thread.sleep(defaultTimeout + increase); - service.checkStates(); - activityVerify(false); - - reset(telemetrySubscriptionService); - - service.onDeviceActivity(tenantId, deviceId, System.currentTimeMillis()); - activityVerify(true); - Thread.sleep(newTimeout); - service.checkStates(); - activityVerify(false); - } + @MethodSource + @ParameterizedTest + void testOnDeviceInactivityTimeoutUpdate(boolean initialActivityStatus, long newInactivityTimeout, boolean expectedActivityStatus) { + // GIVEN + doReturn(200L).when(service).getCurrentTimeMillis(); - @Test - public void increaseSmallInactivityForInactiveDeviceTest() throws Exception { - final long defaultTimeout = 1; - initStateService(defaultTimeout); - DeviceState deviceState = DeviceState.builder().build(); - DeviceStateData deviceStateData = DeviceStateData.builder() - .tenantId(tenantId) - .deviceId(deviceId) - .state(deviceState) - .metaData(new TbMsgMetaData()) + var deviceState = DeviceState.builder() + .active(initialActivityStatus) + .lastActivityTime(100L) .build(); - service.deviceStates.put(deviceId, deviceStateData); - service.getPartitionedEntities(tpi).add(deviceId); - - service.onDeviceActivity(tenantId, deviceId, System.currentTimeMillis()); - activityVerify(true); - Thread.sleep(defaultTimeout); - service.checkStates(); - activityVerify(false); - - reset(telemetrySubscriptionService); - - long newTimeout = 1; - Thread.sleep(newTimeout); - verify(telemetrySubscriptionService, never()).saveAttributes(argThat(request -> - request.getEntityId().equals(deviceId) && request.getEntries().get(0).getKey().equals(ACTIVITY_STATE) - )); - } - - @Test - public void decreaseInactivityForActiveDeviceTest() throws Exception { - final long defaultTimeout = 1000; - initStateService(defaultTimeout); - DeviceState deviceState = DeviceState.builder().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); - - service.onDeviceActivity(tenantId, deviceId, System.currentTimeMillis()); - activityVerify(true); - - long newTimeout = 1; - Thread.sleep(newTimeout); - - service.onDeviceInactivityTimeoutUpdate(tenantId, deviceId, newTimeout); - activityVerify(false); - reset(telemetrySubscriptionService); - - service.onDeviceInactivityTimeoutUpdate(tenantId, deviceId, defaultTimeout); - activityVerify(true); - Thread.sleep(defaultTimeout); - service.checkStates(); - activityVerify(false); - } - - @Test - public void decreaseInactivityForInactiveDeviceTest() throws Exception { - final long defaultTimeout = 1000; - initStateService(defaultTimeout); - DeviceState deviceState = DeviceState.builder().build(); - DeviceStateData deviceStateData = DeviceStateData.builder() + var deviceStateData = DeviceStateData.builder() .tenantId(tenantId) .deviceId(deviceId) .state(deviceState) @@ -679,31 +563,44 @@ public class DefaultDeviceStateServiceTest { service.deviceStates.put(deviceId, deviceStateData); service.getPartitionedEntities(tpi).add(deviceId); - service.onDeviceActivity(tenantId, deviceId, System.currentTimeMillis()); - activityVerify(true); - Thread.sleep(defaultTimeout); - service.checkStates(); - activityVerify(false); - reset(telemetrySubscriptionService); + mockSuccessfulSaveAttributes(); - long newTimeout = 1; + // WHEN + service.onDeviceInactivityTimeoutUpdate(tenantId, deviceId, newInactivityTimeout); - service.onDeviceInactivityTimeoutUpdate(tenantId, deviceId, newTimeout); - verify(telemetrySubscriptionService, never()).saveAttributes(argThat(request -> - request.getEntityId().equals(deviceId) && request.getEntries().get(0).getKey().equals(ACTIVITY_STATE) - )); + // THEN + long expectedInactivityTimeout = newInactivityTimeout != 0 ? newInactivityTimeout : defaultInactivityTimeoutMs; + assertThat(deviceState.getInactivityTimeout()).isEqualTo(expectedInactivityTimeout); + + assertThat(deviceState.isActive()).isEqualTo(expectedActivityStatus); + if (initialActivityStatus != expectedActivityStatus) { + then(telemetrySubscriptionService).should().saveAttributesInternal(argThat(request -> { + AttributeKvEntry entry = request.getEntries().get(0); + return request.getEntityId().equals(deviceId) && entry.getKey().equals(ACTIVITY_STATE) && entry.getValue().equals(expectedActivityStatus); + })); + } } - private void activityVerify(boolean isActive) { - verify(telemetrySubscriptionService).saveAttributes(argThat(request -> - request.getEntityId().equals(deviceId) && - request.getEntries().get(0).getKey().equals(ACTIVITY_STATE) && - request.getEntries().get(0).getValue().equals(isActive) - )); + // to simplify test, these arguments assume that the current time is 200 and the last activity time is 100 + private static Stream testOnDeviceInactivityTimeoutUpdate() { + return Stream.of( + Arguments.of(true, 1L, false), + Arguments.of(true, 50L, false), + Arguments.of(true, 99L, false), + Arguments.of(true, 100L, false), + Arguments.of(true, 101L, true), + Arguments.of(true, 0L, true), // should use default inactivity timeout of 10 minutes + Arguments.of(false, 1L, false), + Arguments.of(false, 50L, false), + Arguments.of(false, 99L, false), + Arguments.of(false, 100L, false), + Arguments.of(false, 101L, true), + Arguments.of(false, 0L, true) // should use default inactivity timeout of 10 minutes + ); } @Test - public void givenStateDataIsNull_whenUpdateActivityState_thenShouldCleanupDevice() { + void givenStateDataIsNull_whenUpdateActivityState_thenShouldCleanupDevice() { // GIVEN service.deviceStates.put(deviceId, deviceStateDataMock); @@ -719,7 +616,7 @@ public class DefaultDeviceStateServiceTest { @ParameterizedTest @MethodSource("provideParametersForUpdateActivityState") - public void givenTestParameters_whenUpdateActivityState_thenShouldBeInTheExpectedStateAndPerformExpectedActions( + void givenTestParameters_whenUpdateActivityState_thenShouldBeInTheExpectedStateAndPerformExpectedActions( boolean activityState, long previousActivityTime, long lastReportedActivity, long inactivityAlarmTime, long expectedInactivityAlarmTime, boolean shouldSetInactivityAlarmTimeToZero, boolean shouldUpdateActivityStateToActive @@ -739,13 +636,15 @@ public class DefaultDeviceStateServiceTest { .metaData(new TbMsgMetaData()) .build(); + mockSuccessfulSaveAttributes(); + // WHEN service.updateActivityState(deviceId, deviceStateData, lastReportedActivity); // THEN assertThat(deviceState.isActive()).isEqualTo(true); assertThat(deviceState.getLastActivityTime()).isEqualTo(lastReportedActivity); - then(telemetrySubscriptionService).should().saveAttributes(argThat(request -> + then(telemetrySubscriptionService).should().saveAttributesInternal(argThat(request -> request.getEntityId().equals(deviceId) && request.getEntries().get(0).getKey().equals(LAST_ACTIVITY_TIME) && request.getEntries().get(0).getValue().equals(lastReportedActivity) @@ -753,7 +652,7 @@ public class DefaultDeviceStateServiceTest { assertThat(deviceState.getLastInactivityAlarmTime()).isEqualTo(expectedInactivityAlarmTime); if (shouldSetInactivityAlarmTimeToZero) { - then(telemetrySubscriptionService).should().saveAttributes(argThat(request -> + then(telemetrySubscriptionService).should().saveAttributesInternal(argThat(request -> request.getEntityId().equals(deviceId) && request.getEntries().get(0).getKey().equals(INACTIVITY_ALARM_TIME) && request.getEntries().get(0).getValue().equals(0L) @@ -761,7 +660,7 @@ public class DefaultDeviceStateServiceTest { } if (shouldUpdateActivityStateToActive) { - then(telemetrySubscriptionService).should().saveAttributes(argThat(request -> + then(telemetrySubscriptionService).should().saveAttributesInternal(argThat(request -> request.getEntityId().equals(deviceId) && request.getEntries().get(0).getKey().equals(ACTIVITY_STATE) && request.getEntries().get(0).getValue().equals(true) @@ -809,59 +708,8 @@ public class DefaultDeviceStateServiceTest { ); } - @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().saveAttributes(argThat(request -> - request.getEntityId().equals(deviceId) && request.getEntries().get(0).getKey().equals(ACTIVITY_STATE) && - request.getEntries().get(0).getValue().equals(false) - )); - } - } - - 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() { + void givenStateDataIsNull_whenUpdateInactivityTimeoutIfExpired_thenShouldCleanupDevice() { // GIVEN service.deviceStates.put(deviceId, deviceStateDataMock); @@ -875,7 +723,7 @@ public class DefaultDeviceStateServiceTest { } @Test - public void givenNotMyPartition_whenUpdateInactivityTimeoutIfExpired_thenShouldCleanupDevice() { + void givenNotMyPartition_whenUpdateInactivityTimeoutIfExpired_thenShouldCleanupDevice() { // GIVEN long currentTime = System.currentTimeMillis(); @@ -911,7 +759,7 @@ public class DefaultDeviceStateServiceTest { @ParameterizedTest @MethodSource("provideParametersForUpdateInactivityStateIfExpired") - public void givenTestParameters_whenUpdateInactivityStateIfExpired_thenShouldBeInTheExpectedStateAndPerformExpectedActions( + void givenTestParameters_whenUpdateInactivityStateIfExpired_thenShouldBeInTheExpectedStateAndPerformExpectedActions( boolean activityState, long ts, long lastActivityTime, long lastInactivityAlarmTime, long inactivityTimeout, long deviceCreationTime, boolean expectedActivityState, long expectedLastInactivityAlarmTime, boolean shouldUpdateActivityStateToInactive ) { @@ -933,6 +781,7 @@ public class DefaultDeviceStateServiceTest { if (shouldUpdateActivityStateToInactive) { given(partitionService.resolve(ServiceType.TB_CORE, tenantId, deviceId)).willReturn(tpi); + mockSuccessfulSaveAttributes(); } // WHEN @@ -943,7 +792,7 @@ public class DefaultDeviceStateServiceTest { assertThat(state.getLastInactivityAlarmTime()).isEqualTo(expectedLastInactivityAlarmTime); if (shouldUpdateActivityStateToInactive) { - then(telemetrySubscriptionService).should().saveAttributes(argThat(request -> + then(telemetrySubscriptionService).should().saveAttributesInternal(argThat(request -> request.getEntityId().equals(deviceId) && request.getEntries().get(0).getKey().equals(ACTIVITY_STATE) && request.getEntries().get(0).getValue().equals(false) )); @@ -961,7 +810,7 @@ public class DefaultDeviceStateServiceTest { assertThat(actualNotification.getDeviceId()).isEqualTo(deviceId); assertThat(actualNotification.isActive()).isFalse(); - then(telemetrySubscriptionService).should().saveAttributes(argThat(request -> + then(telemetrySubscriptionService).should().saveAttributesInternal(argThat(request -> request.getTenantId().equals(tenantId) && request.getEntityId().equals(deviceId) && request.getScope().equals(AttributeScope.SERVER_SCOPE) && request.getEntries().get(0).getKey().equals(INACTIVITY_ALARM_TIME) && @@ -1033,7 +882,80 @@ public class DefaultDeviceStateServiceTest { } @Test - public void givenConcurrentAccess_whenGetOrFetchDeviceStateData_thenFetchDeviceStateDataInvokedOnce() { + void givenInactiveDevice_whenActivityStatusChangesToActiveButFailedToSaveUpdatedActivityStatus_thenShouldNotUpdateCache() { + // GIVEN + doReturn(200L).when(service).getCurrentTimeMillis(); + + var deviceState = DeviceState.builder() + .active(false) + .lastActivityTime(100L) + .inactivityTimeout(50L) + .build(); + + var deviceStateData = DeviceStateData.builder() + .tenantId(tenantId) + .deviceId(deviceId) + .state(deviceState) + .metaData(TbMsgMetaData.EMPTY) + .build(); + + service.deviceStates.put(deviceId, deviceStateData); + service.getPartitionedEntities(tpi).add(deviceId); + + when(telemetrySubscriptionService.saveAttributesInternal(any(AttributesSaveRequest.class))) + .thenAnswer(invocation -> { + AttributesSaveRequest request = invocation.getArgument(0); + AttributeKvEntry entry = request.getEntries().get(0); + return entry.getKey().equals(ACTIVITY_STATE) ? + Futures.immediateFailedFuture(new RuntimeException("failed to save")) : + Futures.immediateFuture(generateRandomVersions(1)); + }); + + // WHEN + service.onDeviceActivity(tenantId, deviceId, 220L); + + // THEN + assertThat(deviceState.isActive()).isFalse(); + } + + @Test + void givenActiveDevice_whenActivityStatusChangesToInactiveButFailedToSaveUpdatedActivityStatus_thenShouldNotUpdateCache() { + // GIVEN + var deviceState = DeviceState.builder() + .active(true) + .lastActivityTime(100L) + .inactivityTimeout(50L) + .build(); + + var deviceStateData = DeviceStateData.builder() + .tenantId(tenantId) + .deviceId(deviceId) + .state(deviceState) + .metaData(TbMsgMetaData.EMPTY) + .build(); + + service.deviceStates.put(deviceId, deviceStateData); + service.getPartitionedEntities(tpi).add(deviceId); + + when(telemetrySubscriptionService.saveAttributesInternal(any(AttributesSaveRequest.class))) + .thenAnswer(invocation -> { + AttributesSaveRequest request = invocation.getArgument(0); + AttributeKvEntry entry = request.getEntries().get(0); + return entry.getKey().equals(ACTIVITY_STATE) ? + Futures.immediateFailedFuture(new RuntimeException("failed to save")) : + Futures.immediateFuture(generateRandomVersions(1)); + }); + + // WHEN + doReturn(200L).when(service).getCurrentTimeMillis(); + service.checkStates(); + + // THEN + assertThat(deviceState.isActive()).isTrue(); + } + + @Test + void givenConcurrentAccess_whenGetOrFetchDeviceStateData_thenFetchDeviceStateDataInvokedOnce() { doAnswer(invocation -> { Thread.sleep(100); return deviceStateDataMock; @@ -1069,10 +991,8 @@ public class DefaultDeviceStateServiceTest { } @Test - public void givenDeviceAdded_whenOnQueueMsg_thenShouldCacheAndSaveActivityToFalse() throws InterruptedException { + void givenDeviceAdded_whenOnQueueMsg_thenShouldCacheAndSaveActivityToFalse() { // GIVEN - final long defaultTimeout = 1000; - initStateService(defaultTimeout); given(deviceService.findDeviceById(any(TenantId.class), any(DeviceId.class))).willReturn(new Device(deviceId)); given(attributesService.find(any(TenantId.class), any(EntityId.class), any(AttributeScope.class), anyCollection())).willReturn(Futures.immediateFuture(Collections.emptyList())); @@ -1086,13 +1006,15 @@ public class DefaultDeviceStateServiceTest { .setDeleted(false) .build(); + mockSuccessfulSaveAttributes(); + // WHEN service.onQueueMsg(proto, TbCallback.EMPTY); // THEN await().atMost(1, TimeUnit.SECONDS).untilAsserted(() -> { assertThat(service.deviceStates.get(deviceId).getState().isActive()).isEqualTo(false); - then(telemetrySubscriptionService).should().saveAttributes(argThat(request -> + then(telemetrySubscriptionService).should().saveAttributesInternal(argThat(request -> request.getEntityId().equals(deviceId) && request.getEntries().get(0).getKey().equals(ACTIVITY_STATE) && request.getEntries().get(0).getValue().equals(false) )); @@ -1100,14 +1022,12 @@ public class DefaultDeviceStateServiceTest { } @Test - public void givenDeviceActivityEventHappenedAfterAdded_whenOnDeviceActivity_thenShouldCacheAndSaveActivityToTrue() throws InterruptedException { + void givenDeviceActivityEventHappenedAfterAdded_whenOnDeviceActivity_thenShouldCacheAndSaveActivityToTrue() { // GIVEN - final long defaultTimeout = 1000; - initStateService(defaultTimeout); long currentTime = System.currentTimeMillis(); DeviceState deviceState = DeviceState.builder() .active(false) - .inactivityTimeout(service.getDefaultInactivityTimeoutInSec()) + .inactivityTimeout(defaultInactivityTimeoutMs) .build(); DeviceStateData stateData = DeviceStateData.builder() .tenantId(tenantId) @@ -1118,12 +1038,14 @@ public class DefaultDeviceStateServiceTest { .build(); service.deviceStates.put(deviceId, stateData); + mockSuccessfulSaveAttributes(); + // WHEN service.onDeviceActivity(tenantId, deviceId, currentTime); // THEN ArgumentCaptor attributeRequestCaptor = ArgumentCaptor.forClass(AttributesSaveRequest.class); - then(telemetrySubscriptionService).should(times(2)).saveAttributes(attributeRequestCaptor.capture()); + then(telemetrySubscriptionService).should(times(2)).saveAttributesInternal(attributeRequestCaptor.capture()); await().atMost(1, TimeUnit.SECONDS).untilAsserted(() -> { assertThat(service.deviceStates.get(deviceId).getState().isActive()).isEqualTo(true); @@ -1151,15 +1073,14 @@ public class DefaultDeviceStateServiceTest { } @Test - public void givenDeviceActivityEventHappenedBeforeAdded_whenOnQueueMsg_thenShouldSaveActivityStateUsingValueFromCache() throws InterruptedException { + void givenDeviceActivityEventHappenedBeforeAdded_whenOnQueueMsg_thenShouldSaveActivityStateUsingValueFromCache() { // GIVEN - final long defaultTimeout = 1000; - initStateService(defaultTimeout); given(deviceService.findDeviceById(any(TenantId.class), any(DeviceId.class))).willReturn(new Device(deviceId)); given(attributesService.find(any(TenantId.class), any(EntityId.class), any(AttributeScope.class), anyCollection())).willReturn(Futures.immediateFuture(Collections.emptyList())); long currentTime = System.currentTimeMillis(); - DeviceState deviceState = DeviceState.builder() + + var deviceState = DeviceState.builder() .active(true) .lastConnectTime(currentTime - 8000) .lastActivityTime(currentTime - 4000) @@ -1167,16 +1088,20 @@ public class DefaultDeviceStateServiceTest { .lastInactivityAlarmTime(0) .inactivityTimeout(3000) .build(); - DeviceStateData stateData = DeviceStateData.builder() + + var stateData = DeviceStateData.builder() .tenantId(tenantId) .deviceId(deviceId) .deviceCreationTime(currentTime - 10000) .state(deviceState) .build(); + service.deviceStates.put(deviceId, stateData); + mockSuccessfulSaveAttributes(); + // WHEN - TransportProtos.DeviceStateServiceMsgProto proto = TransportProtos.DeviceStateServiceMsgProto.newBuilder() + var proto = TransportProtos.DeviceStateServiceMsgProto.newBuilder() .setTenantIdMSB(tenantId.getId().getMostSignificantBits()) .setTenantIdLSB(tenantId.getId().getLeastSignificantBits()) .setDeviceIdMSB(deviceId.getId().getMostSignificantBits()) @@ -1190,11 +1115,25 @@ public class DefaultDeviceStateServiceTest { // THEN await().atMost(1, TimeUnit.SECONDS).untilAsserted(() -> { assertThat(service.deviceStates.get(deviceId).getState().isActive()).isEqualTo(true); - then(telemetrySubscriptionService).should().saveAttributes(argThat(request -> + then(telemetrySubscriptionService).should().saveAttributesInternal(argThat(request -> request.getEntityId().equals(deviceId) && request.getEntries().get(0).getKey().equals(ACTIVITY_STATE) && request.getEntries().get(0).getValue().equals(true) )); }); } + private void mockSuccessfulSaveAttributes() { + lenient().when(telemetrySubscriptionService.saveAttributesInternal(any())).thenAnswer(invocation -> { + AttributesSaveRequest request = invocation.getArgument(0); + return Futures.immediateFuture(generateRandomVersions(request.getEntries().size())); + }); + } + + private static List generateRandomVersions(int n) { + return ThreadLocalRandom.current() + .longs(n) + .boxed() + .toList(); + } + }