From b56f75a99675a110b6efe23fd17162c02e108b98 Mon Sep 17 00:00:00 2001 From: IrynaMatveieva Date: Thu, 29 Aug 2024 10:45:34 +0300 Subject: [PATCH 1/4] added validation on init --- .../engine/aws/lambda/TbAwsLambdaNode.java | 25 ++++++- .../aws/lambda/TbAwsLambdaNodeTest.java | 72 +++++++++++++++++-- 2 files changed, 90 insertions(+), 7 deletions(-) diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/aws/lambda/TbAwsLambdaNode.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/aws/lambda/TbAwsLambdaNode.java index 2512125ff1..d20ff2f75c 100644 --- a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/aws/lambda/TbAwsLambdaNode.java +++ b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/aws/lambda/TbAwsLambdaNode.java @@ -62,9 +62,7 @@ public class TbAwsLambdaNode extends TbAbstractExternalNode { @Override public void init(TbContext ctx, TbNodeConfiguration configuration) throws TbNodeException { config = TbNodeUtils.convert(configuration, TbAwsLambdaNodeConfiguration.class); - if (StringUtils.isBlank(config.getFunctionName())) { - throw new TbNodeException("Function name must be set!", true); - } + validateConfig(); try { AWSCredentials awsCredentials = new BasicAWSCredentials(config.getAccessKey(), config.getSecretKey()); client = AWSLambdaAsyncClientBuilder.standard() @@ -138,6 +136,27 @@ public class TbAwsLambdaNode extends TbAbstractExternalNode { return TbMsg.transformMsgMetadata(origMsg, metaData); } + private void validateConfig() throws TbNodeException { + if (StringUtils.isBlank(config.getFunctionName())) { + throw new TbNodeException("Function name must be set!", true); + } + if (StringUtils.isBlank(config.getAccessKey())) { + throw new TbNodeException("Access Key must be set!", true); + } + if (StringUtils.isBlank(config.getSecretKey())) { + throw new TbNodeException("Secret Access Key must be set!", true); + } + if (StringUtils.isBlank(config.getRegion())) { + throw new TbNodeException("Region must be set!", true); + } + if (config.getConnectionTimeout() < 0) { + throw new TbNodeException("Min connection timeout is 0!", true); + } + if (config.getRequestTimeout() < 0) { + throw new TbNodeException("Min request timeout is 0!", true); + } + } + @Override public void destroy() { if (client != null) { diff --git a/rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/aws/lambda/TbAwsLambdaNodeTest.java b/rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/aws/lambda/TbAwsLambdaNodeTest.java index 4ee954991c..eed3613a48 100644 --- a/rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/aws/lambda/TbAwsLambdaNodeTest.java +++ b/rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/aws/lambda/TbAwsLambdaNodeTest.java @@ -75,10 +75,14 @@ public class TbAwsLambdaNodeTest { void setUp() { node = new TbAwsLambdaNode(); config = new TbAwsLambdaNodeConfiguration().defaultConfiguration(); + config.setAccessKey("accessKey"); + config.setSecretKey("secretKey"); + config.setFunctionName("new-function"); } @Test public void verifyDefaultConfig() { + config = new TbAwsLambdaNodeConfiguration().defaultConfiguration(); assertThat(config.getAccessKey()).isNull(); assertThat(config.getSecretKey()).isNull(); assertThat(config.getRegion()).isEqualTo(("us-east-1")); @@ -97,7 +101,70 @@ public class TbAwsLambdaNodeTest { var configuration = new TbNodeConfiguration(JacksonUtil.valueToTree(config)); assertThatThrownBy(() -> node.init(ctx, configuration)) .isInstanceOf(TbNodeException.class) - .hasMessage("Function name must be set!"); + .hasMessage("Function name must be set!") + .extracting(e -> ((TbNodeException) e).isUnrecoverable()) + .isEqualTo(true); + } + + @ParameterizedTest + @NullAndEmptySource + @ValueSource(strings = " ") + public void givenInvalidAccessKey_whenInit_thenThrowsException(String accessKey) { + config.setAccessKey(accessKey); + var configuration = new TbNodeConfiguration(JacksonUtil.valueToTree(config)); + assertThatThrownBy(() -> node.init(ctx, configuration)) + .isInstanceOf(TbNodeException.class) + .hasMessage("Access Key must be set!") + .extracting(e -> ((TbNodeException) e).isUnrecoverable()) + .isEqualTo(true); + } + + @ParameterizedTest + @NullAndEmptySource + @ValueSource(strings = " ") + public void givenInvalidSecretAccessKey_whenInit_thenThrowsException(String secretAccessKey) { + config.setSecretKey(secretAccessKey); + var configuration = new TbNodeConfiguration(JacksonUtil.valueToTree(config)); + assertThatThrownBy(() -> node.init(ctx, configuration)) + .isInstanceOf(TbNodeException.class) + .hasMessage("Secret Access Key must be set!") + .extracting(e -> ((TbNodeException) e).isUnrecoverable()) + .isEqualTo(true); + } + + @ParameterizedTest + @NullAndEmptySource + @ValueSource(strings = " ") + public void givenInvalidRegion_whenInit_thenThrowsException(String region) { + config.setRegion(region); + var configuration = new TbNodeConfiguration(JacksonUtil.valueToTree(config)); + assertThatThrownBy(() -> node.init(ctx, configuration)) + .isInstanceOf(TbNodeException.class) + .hasMessage("Region must be set!") + .extracting(e -> ((TbNodeException) e).isUnrecoverable()) + .isEqualTo(true); + } + + @Test + public void givenInvalidConnectionTimeout_whenInit_thenThrowsException() { + config.setConnectionTimeout(-100); + var configuration = new TbNodeConfiguration(JacksonUtil.valueToTree(config)); + assertThatThrownBy(() -> node.init(ctx, configuration)) + .isInstanceOf(TbNodeException.class) + .hasMessage("Min connection timeout is 0!") + .extracting(e -> ((TbNodeException) e).isUnrecoverable()) + .isEqualTo(true); + } + + @Test + public void givenInvalidRequestTimeout_whenInit_thenThrowsException() { + config.setRequestTimeout(-100); + var configuration = new TbNodeConfiguration(JacksonUtil.valueToTree(config)); + assertThatThrownBy(() -> node.init(ctx, configuration)) + .isInstanceOf(TbNodeException.class) + .hasMessage("Min request timeout is 0!") + .extracting(e -> ((TbNodeException) e).isUnrecoverable()) + .isEqualTo(true); } @ParameterizedTest @@ -281,9 +348,6 @@ public class TbAwsLambdaNodeTest { } private void init() { - config.setAccessKey("accessKey"); - config.setSecretKey("secretKey"); - config.setFunctionName("new-function"); ReflectionTestUtils.setField(node, "client", clientMock); ReflectionTestUtils.setField(node, "config", config); } From 5256cae80406a06842cff423a2ecdfe2361fe2b0 Mon Sep 17 00:00:00 2001 From: IrynaMatveieva Date: Thu, 29 Aug 2024 12:35:30 +0300 Subject: [PATCH 2/4] used jakarta validation constraints --- .../engine/aws/lambda/TbAwsLambdaNode.java | 26 ++------ .../lambda/TbAwsLambdaNodeConfiguration.java | 8 +++ .../aws/lambda/TbAwsLambdaNodeTest.java | 59 +++++++------------ 3 files changed, 32 insertions(+), 61 deletions(-) diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/aws/lambda/TbAwsLambdaNode.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/aws/lambda/TbAwsLambdaNode.java index d20ff2f75c..06dc0310ae 100644 --- a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/aws/lambda/TbAwsLambdaNode.java +++ b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/aws/lambda/TbAwsLambdaNode.java @@ -39,6 +39,8 @@ import org.thingsboard.server.common.msg.TbMsgMetaData; import java.nio.ByteBuffer; import java.util.concurrent.TimeUnit; +import static org.thingsboard.server.dao.service.ConstraintValidator.validateFields; + @Slf4j @RuleNode( type = ComponentType.EXTERNAL, @@ -62,7 +64,8 @@ public class TbAwsLambdaNode extends TbAbstractExternalNode { @Override public void init(TbContext ctx, TbNodeConfiguration configuration) throws TbNodeException { config = TbNodeUtils.convert(configuration, TbAwsLambdaNodeConfiguration.class); - validateConfig(); + String errorPrefix = "'" + ctx.getSelf().getName() + "' node configuration is invalid: "; + validateFields(config, errorPrefix); try { AWSCredentials awsCredentials = new BasicAWSCredentials(config.getAccessKey(), config.getSecretKey()); client = AWSLambdaAsyncClientBuilder.standard() @@ -136,27 +139,6 @@ public class TbAwsLambdaNode extends TbAbstractExternalNode { return TbMsg.transformMsgMetadata(origMsg, metaData); } - private void validateConfig() throws TbNodeException { - if (StringUtils.isBlank(config.getFunctionName())) { - throw new TbNodeException("Function name must be set!", true); - } - if (StringUtils.isBlank(config.getAccessKey())) { - throw new TbNodeException("Access Key must be set!", true); - } - if (StringUtils.isBlank(config.getSecretKey())) { - throw new TbNodeException("Secret Access Key must be set!", true); - } - if (StringUtils.isBlank(config.getRegion())) { - throw new TbNodeException("Region must be set!", true); - } - if (config.getConnectionTimeout() < 0) { - throw new TbNodeException("Min connection timeout is 0!", true); - } - if (config.getRequestTimeout() < 0) { - throw new TbNodeException("Min request timeout is 0!", true); - } - } - @Override public void destroy() { if (client != null) { diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/aws/lambda/TbAwsLambdaNodeConfiguration.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/aws/lambda/TbAwsLambdaNodeConfiguration.java index c7c66599d5..b395f6b15e 100644 --- a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/aws/lambda/TbAwsLambdaNodeConfiguration.java +++ b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/aws/lambda/TbAwsLambdaNodeConfiguration.java @@ -15,6 +15,8 @@ */ package org.thingsboard.rule.engine.aws.lambda; +import jakarta.validation.constraints.Min; +import jakarta.validation.constraints.NotBlank; import lombok.Data; import org.thingsboard.rule.engine.api.NodeConfiguration; @@ -23,12 +25,18 @@ public class TbAwsLambdaNodeConfiguration implements NodeConfiguration node.init(ctx, configuration)) - .isInstanceOf(TbNodeException.class) - .hasMessage("Function name must be set!") - .extracting(e -> ((TbNodeException) e).isUnrecoverable()) - .isEqualTo(true); + verifyDataValidationExceptionOnInit(); } @ParameterizedTest @@ -111,12 +107,7 @@ public class TbAwsLambdaNodeTest { @ValueSource(strings = " ") public void givenInvalidAccessKey_whenInit_thenThrowsException(String accessKey) { config.setAccessKey(accessKey); - var configuration = new TbNodeConfiguration(JacksonUtil.valueToTree(config)); - assertThatThrownBy(() -> node.init(ctx, configuration)) - .isInstanceOf(TbNodeException.class) - .hasMessage("Access Key must be set!") - .extracting(e -> ((TbNodeException) e).isUnrecoverable()) - .isEqualTo(true); + verifyDataValidationExceptionOnInit(); } @ParameterizedTest @@ -124,12 +115,7 @@ public class TbAwsLambdaNodeTest { @ValueSource(strings = " ") public void givenInvalidSecretAccessKey_whenInit_thenThrowsException(String secretAccessKey) { config.setSecretKey(secretAccessKey); - var configuration = new TbNodeConfiguration(JacksonUtil.valueToTree(config)); - assertThatThrownBy(() -> node.init(ctx, configuration)) - .isInstanceOf(TbNodeException.class) - .hasMessage("Secret Access Key must be set!") - .extracting(e -> ((TbNodeException) e).isUnrecoverable()) - .isEqualTo(true); + verifyDataValidationExceptionOnInit(); } @ParameterizedTest @@ -137,34 +123,19 @@ public class TbAwsLambdaNodeTest { @ValueSource(strings = " ") public void givenInvalidRegion_whenInit_thenThrowsException(String region) { config.setRegion(region); - var configuration = new TbNodeConfiguration(JacksonUtil.valueToTree(config)); - assertThatThrownBy(() -> node.init(ctx, configuration)) - .isInstanceOf(TbNodeException.class) - .hasMessage("Region must be set!") - .extracting(e -> ((TbNodeException) e).isUnrecoverable()) - .isEqualTo(true); + verifyDataValidationExceptionOnInit(); } @Test public void givenInvalidConnectionTimeout_whenInit_thenThrowsException() { config.setConnectionTimeout(-100); - var configuration = new TbNodeConfiguration(JacksonUtil.valueToTree(config)); - assertThatThrownBy(() -> node.init(ctx, configuration)) - .isInstanceOf(TbNodeException.class) - .hasMessage("Min connection timeout is 0!") - .extracting(e -> ((TbNodeException) e).isUnrecoverable()) - .isEqualTo(true); + verifyDataValidationExceptionOnInit(); } @Test public void givenInvalidRequestTimeout_whenInit_thenThrowsException() { config.setRequestTimeout(-100); - var configuration = new TbNodeConfiguration(JacksonUtil.valueToTree(config)); - assertThatThrownBy(() -> node.init(ctx, configuration)) - .isInstanceOf(TbNodeException.class) - .hasMessage("Min request timeout is 0!") - .extracting(e -> ((TbNodeException) e).isUnrecoverable()) - .isEqualTo(true); + verifyDataValidationExceptionOnInit(); } @ParameterizedTest @@ -347,6 +318,16 @@ public class TbAwsLambdaNodeTest { assertThat(throwableCaptor.getValue()).isInstanceOf(AWSLambdaException.class).hasMessageStartingWith(errorMsg); } + private void verifyDataValidationExceptionOnInit() { + RuleNode ruleNode = new RuleNode(); + ruleNode.setName("test"); + when(ctx.getSelf()).thenReturn(ruleNode); + String errorPrefix = "'test' node configuration is invalid: "; + assertThatThrownBy(() -> node.init(ctx, new TbNodeConfiguration(JacksonUtil.valueToTree(config)))) + .isInstanceOf(DataValidationException.class) + .hasMessageContaining(errorPrefix); + } + private void init() { ReflectionTestUtils.setField(node, "client", clientMock); ReflectionTestUtils.setField(node, "config", config); From 7a84dc655dbff28805e54061d14275f375c884f8 Mon Sep 17 00:00:00 2001 From: IrynaMatveieva Date: Thu, 5 Sep 2024 12:23:43 +0300 Subject: [PATCH 3/4] added catch for DataValidationException --- .../engine/aws/lambda/TbAwsLambdaNode.java | 5 ++++- .../aws/lambda/TbAwsLambdaNodeTest.java | 22 ++++++++++--------- 2 files changed, 16 insertions(+), 11 deletions(-) diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/aws/lambda/TbAwsLambdaNode.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/aws/lambda/TbAwsLambdaNode.java index 06dc0310ae..1dbc81bfa1 100644 --- a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/aws/lambda/TbAwsLambdaNode.java +++ b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/aws/lambda/TbAwsLambdaNode.java @@ -35,6 +35,7 @@ import org.thingsboard.server.common.data.StringUtils; import org.thingsboard.server.common.data.plugin.ComponentType; import org.thingsboard.server.common.msg.TbMsg; import org.thingsboard.server.common.msg.TbMsgMetaData; +import org.thingsboard.server.dao.exception.DataValidationException; import java.nio.ByteBuffer; import java.util.concurrent.TimeUnit; @@ -65,8 +66,8 @@ public class TbAwsLambdaNode extends TbAbstractExternalNode { public void init(TbContext ctx, TbNodeConfiguration configuration) throws TbNodeException { config = TbNodeUtils.convert(configuration, TbAwsLambdaNodeConfiguration.class); String errorPrefix = "'" + ctx.getSelf().getName() + "' node configuration is invalid: "; - validateFields(config, errorPrefix); try { + validateFields(config, errorPrefix); AWSCredentials awsCredentials = new BasicAWSCredentials(config.getAccessKey(), config.getSecretKey()); client = AWSLambdaAsyncClientBuilder.standard() .withCredentials(new AWSStaticCredentialsProvider(awsCredentials)) @@ -75,6 +76,8 @@ public class TbAwsLambdaNode extends TbAbstractExternalNode { .withConnectionTimeout((int) TimeUnit.SECONDS.toMillis(config.getConnectionTimeout())) .withRequestTimeout((int) TimeUnit.SECONDS.toMillis(config.getRequestTimeout()))) .build(); + } catch (DataValidationException e) { + throw new TbNodeException(e, true); } catch (Exception e) { throw new TbNodeException(e); } diff --git a/rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/aws/lambda/TbAwsLambdaNodeTest.java b/rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/aws/lambda/TbAwsLambdaNodeTest.java index d730b9b837..14ae9bef20 100644 --- a/rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/aws/lambda/TbAwsLambdaNodeTest.java +++ b/rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/aws/lambda/TbAwsLambdaNodeTest.java @@ -36,6 +36,7 @@ import org.springframework.test.util.ReflectionTestUtils; import org.thingsboard.common.util.JacksonUtil; import org.thingsboard.rule.engine.api.TbContext; import org.thingsboard.rule.engine.api.TbNodeConfiguration; +import org.thingsboard.rule.engine.api.TbNodeException; import org.thingsboard.rule.engine.api.util.TbNodeUtils; import org.thingsboard.server.common.data.StringUtils; import org.thingsboard.server.common.data.id.DeviceId; @@ -43,7 +44,6 @@ import org.thingsboard.server.common.data.msg.TbMsgType; import org.thingsboard.server.common.data.rule.RuleNode; import org.thingsboard.server.common.msg.TbMsg; import org.thingsboard.server.common.msg.TbMsgMetaData; -import org.thingsboard.server.dao.exception.DataValidationException; import java.nio.ByteBuffer; import java.nio.charset.StandardCharsets; @@ -99,7 +99,7 @@ public class TbAwsLambdaNodeTest { @ValueSource(strings = " ") public void givenInvalidFunctionName_whenInit_thenThrowsException(String funcName) { config.setFunctionName(funcName); - verifyDataValidationExceptionOnInit(); + verifyValidationExceptionOnInit(); } @ParameterizedTest @@ -107,7 +107,7 @@ public class TbAwsLambdaNodeTest { @ValueSource(strings = " ") public void givenInvalidAccessKey_whenInit_thenThrowsException(String accessKey) { config.setAccessKey(accessKey); - verifyDataValidationExceptionOnInit(); + verifyValidationExceptionOnInit(); } @ParameterizedTest @@ -115,7 +115,7 @@ public class TbAwsLambdaNodeTest { @ValueSource(strings = " ") public void givenInvalidSecretAccessKey_whenInit_thenThrowsException(String secretAccessKey) { config.setSecretKey(secretAccessKey); - verifyDataValidationExceptionOnInit(); + verifyValidationExceptionOnInit(); } @ParameterizedTest @@ -123,19 +123,19 @@ public class TbAwsLambdaNodeTest { @ValueSource(strings = " ") public void givenInvalidRegion_whenInit_thenThrowsException(String region) { config.setRegion(region); - verifyDataValidationExceptionOnInit(); + verifyValidationExceptionOnInit(); } @Test public void givenInvalidConnectionTimeout_whenInit_thenThrowsException() { config.setConnectionTimeout(-100); - verifyDataValidationExceptionOnInit(); + verifyValidationExceptionOnInit(); } @Test public void givenInvalidRequestTimeout_whenInit_thenThrowsException() { config.setRequestTimeout(-100); - verifyDataValidationExceptionOnInit(); + verifyValidationExceptionOnInit(); } @ParameterizedTest @@ -318,14 +318,16 @@ public class TbAwsLambdaNodeTest { assertThat(throwableCaptor.getValue()).isInstanceOf(AWSLambdaException.class).hasMessageStartingWith(errorMsg); } - private void verifyDataValidationExceptionOnInit() { + private void verifyValidationExceptionOnInit() { RuleNode ruleNode = new RuleNode(); ruleNode.setName("test"); when(ctx.getSelf()).thenReturn(ruleNode); String errorPrefix = "'test' node configuration is invalid: "; assertThatThrownBy(() -> node.init(ctx, new TbNodeConfiguration(JacksonUtil.valueToTree(config)))) - .isInstanceOf(DataValidationException.class) - .hasMessageContaining(errorPrefix); + .isInstanceOf(TbNodeException.class) + .hasMessageContaining(errorPrefix) + .extracting(e -> ((TbNodeException) e).isUnrecoverable()) + .isEqualTo(true); } private void init() { From 6f2801db0aab0b901524196cc2040f3ff6b3d955 Mon Sep 17 00:00:00 2001 From: Iryna Matveieva <101514424+irynamatveieva@users.noreply.github.com> Date: Thu, 5 Sep 2024 12:38:52 +0300 Subject: [PATCH 4/4] Fixed incorrect display of device state (#11536) * changed logic to save to db activity value from cache * changed check if partition belongs --- .../state/DefaultDeviceStateService.java | 63 +++++----- .../state/DefaultDeviceStateServiceTest.java | 108 ++++++++++++++++++ 2 files changed, 141 insertions(+), 30 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 1c53920a93..6025e874b9 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 @@ -157,7 +157,8 @@ public class DefaultDeviceStateService extends AbstractPartitionBasedService() { @Override - public void onSuccess(@Nullable DeviceStateData state) { + public void onSuccess(DeviceStateData state) { TopicPartitionInfo tpi = partitionService.resolve(ServiceType.TB_CORE, tenantId, device.getId()); - if (addDeviceUsingState(tpi, state)) { - save(deviceId, ACTIVITY_STATE, false); + Set deviceIds = partitionedEntities.get(tpi); + boolean isMyPartition = deviceIds != null; + if (isMyPartition) { + deviceIds.add(state.getDeviceId()); + initializeActivityState(deviceId, state); callback.onSuccess(); } else { - log.debug("[{}][{}] Device belongs to external partition. Probably rebalancing is in progress. Topic: {}" - , tenantId, deviceId, tpi.getFullTopicName()); + log.debug("[{}][{}] Device belongs to external partition. Probably rebalancing is in progress. Topic: {}", tenantId, deviceId, tpi.getFullTopicName()); callback.onFailure(new RuntimeException("Device belongs to external partition " + tpi.getFullTopicName() + "!")); } } @@ -400,6 +403,21 @@ public class DefaultDeviceStateService extends AbstractPartitionBasedService deviceIdSet = partitionedEntities.get(tpi); + if (deviceIdSet != null) { + deviceIdSet.remove(deviceId); + } + } + + private void initializeActivityState(DeviceId deviceId, DeviceStateData fetchedState) { + DeviceStateData cachedState = deviceStates.putIfAbsent(fetchedState.getDeviceId(), fetchedState); + boolean activityState = Objects.requireNonNullElse(cachedState, fetchedState).getState().isActive(); + save(deviceId, ACTIVITY_STATE, activityState); + } + @Override protected Map>> onAddedPartitions(Set addedPartitions) { var result = new HashMap>>(); @@ -436,10 +454,16 @@ public class DefaultDeviceStateService extends AbstractPartitionBasedService deviceIds = partitionedEntities.get(tpi); + boolean isMyPartition = deviceIds != null; + if (isMyPartition) { + deviceIds.add(state.getDeviceId()); + deviceStates.putIfAbsent(state.getDeviceId(), state); + checkAndUpdateState(state.getDeviceId(), state); + } else { + log.debug("[{}] Device belongs to external partition {}", state.getDeviceId(), tpi.getFullTopicName()); } - checkAndUpdateState(state.getDeviceId(), state); } log.info("[{}] Initialized {} out of {} device states", entry.getKey().getPartition().orElse(0), counter.addAndGet(states.size()), entry.getValue().size()); } @@ -475,18 +499,6 @@ public class DefaultDeviceStateService extends AbstractPartitionBasedService deviceIds = partitionedEntities.get(tpi); - if (deviceIds != null) { - deviceIds.add(state.getDeviceId()); - deviceStates.putIfAbsent(state.getDeviceId(), state); - return true; - } else { - log.debug("[{}] Device belongs to external partition {}", state.getDeviceId(), tpi.getFullTopicName()); - return false; - } - } - void checkStates() { try { final long ts = getCurrentTimeMillis(); @@ -619,15 +631,6 @@ public class DefaultDeviceStateService extends AbstractPartitionBasedService deviceIdSet = partitionedEntities.get(tpi); - if (deviceIdSet != null) { - deviceIdSet.remove(deviceId); - } - } - @Override protected void cleanupEntityOnPartitionRemoval(DeviceId deviceId) { cleanupEntity(deviceId); 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 9833396a0a..296191d61b 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,6 +15,7 @@ */ package org.thingsboard.server.service.state; +import com.google.common.util.concurrent.Futures; import org.junit.jupiter.api.BeforeEach; import org.junit.jupiter.api.Test; import org.junit.jupiter.api.extension.ExtendWith; @@ -28,8 +29,10 @@ import org.mockito.junit.jupiter.MockitoExtension; import org.springframework.test.util.ReflectionTestUtils; 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.msg.TbMsgType; import org.thingsboard.server.common.data.notification.rule.trigger.DeviceActivityTrigger; @@ -41,11 +44,13 @@ 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; +import org.thingsboard.server.common.msg.queue.TbCallback; import org.thingsboard.server.common.msg.queue.TopicPartitionInfo; import org.thingsboard.server.dao.attributes.AttributesService; import org.thingsboard.server.dao.device.DeviceService; 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; @@ -66,6 +71,7 @@ import java.util.stream.Stream; import static org.assertj.core.api.Assertions.assertThat; import static org.awaitility.Awaitility.await; import static org.mockito.ArgumentMatchers.any; +import static org.mockito.ArgumentMatchers.anyList; import static org.mockito.ArgumentMatchers.anyLong; import static org.mockito.ArgumentMatchers.eq; import static org.mockito.BDDMockito.given; @@ -1070,4 +1076,106 @@ public class DefaultDeviceStateServiceTest { then(service).should().fetchDeviceStateDataUsingSeparateRequests(deviceId); } + @Test + public void givenDeviceAdded_whenOnQueueMsg_thenShouldCacheAndSaveActivityToFalse() throws InterruptedException { + // 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), anyList())).willReturn(Futures.immediateFuture(Collections.emptyList())); + + TransportProtos.DeviceStateServiceMsgProto proto = TransportProtos.DeviceStateServiceMsgProto.newBuilder() + .setTenantIdMSB(tenantId.getId().getMostSignificantBits()) + .setTenantIdLSB(tenantId.getId().getLeastSignificantBits()) + .setDeviceIdMSB(deviceId.getId().getMostSignificantBits()) + .setDeviceIdLSB(deviceId.getId().getLeastSignificantBits()) + .setAdded(true) + .setUpdated(false) + .setDeleted(false) + .build(); + + // 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().saveAttrAndNotify(eq(TenantId.SYS_TENANT_ID), eq(deviceId), eq(AttributeScope.SERVER_SCOPE), eq(ACTIVITY_STATE), eq(false), any()); + }); + } + + @Test + public void givenDeviceActivityEventHappenedAfterAdded_whenOnDeviceActivity_thenShouldCacheAndSaveActivityToTrue() throws InterruptedException { + // GIVEN + final long defaultTimeout = 1000; + initStateService(defaultTimeout); + long currentTime = System.currentTimeMillis(); + DeviceState deviceState = DeviceState.builder() + .active(false) + .inactivityTimeout(service.getDefaultInactivityTimeoutInSec()) + .build(); + DeviceStateData stateData = DeviceStateData.builder() + .tenantId(tenantId) + .deviceId(deviceId) + .deviceCreationTime(currentTime - 10000) + .state(deviceState) + .metaData(TbMsgMetaData.EMPTY) + .build(); + service.deviceStates.put(deviceId, stateData); + + // WHEN + service.onDeviceActivity(tenantId, deviceId, currentTime); + + // THEN + await().atMost(1, TimeUnit.SECONDS).untilAsserted(() -> { + assertThat(service.deviceStates.get(deviceId).getState().isActive()).isEqualTo(true); + then(telemetrySubscriptionService).should().saveAttrAndNotify(eq(TenantId.SYS_TENANT_ID), eq(deviceId), eq(AttributeScope.SERVER_SCOPE), eq(LAST_ACTIVITY_TIME), eq(currentTime), any()); + then(telemetrySubscriptionService).should().saveAttrAndNotify(eq(TenantId.SYS_TENANT_ID), eq(deviceId), eq(AttributeScope.SERVER_SCOPE), eq(ACTIVITY_STATE), eq(true), any()); + }); + } + + @Test + public void givenDeviceActivityEventHappenedBeforeAdded_whenOnQueueMsg_thenShouldSaveActivityStateUsingValueFromCache() throws InterruptedException { + // 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), anyList())).willReturn(Futures.immediateFuture(Collections.emptyList())); + + 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); + + // WHEN + TransportProtos.DeviceStateServiceMsgProto proto = TransportProtos.DeviceStateServiceMsgProto.newBuilder() + .setTenantIdMSB(tenantId.getId().getMostSignificantBits()) + .setTenantIdLSB(tenantId.getId().getLeastSignificantBits()) + .setDeviceIdMSB(deviceId.getId().getMostSignificantBits()) + .setDeviceIdLSB(deviceId.getId().getLeastSignificantBits()) + .setAdded(true) + .setUpdated(false) + .setDeleted(false) + .build(); + service.onQueueMsg(proto, TbCallback.EMPTY); + + // THEN + await().atMost(1, TimeUnit.SECONDS).untilAsserted(() -> { + assertThat(service.deviceStates.get(deviceId).getState().isActive()).isEqualTo(true); + then(telemetrySubscriptionService).should().saveAttrAndNotify(eq(TenantId.SYS_TENANT_ID), eq(deviceId), eq(AttributeScope.SERVER_SCOPE), eq(ACTIVITY_STATE), eq(true), any()); + }); + } + }