|
|
@ -1,12 +1,12 @@ |
|
|
/** |
|
|
/** |
|
|
* Copyright © 2016-2023 The Thingsboard Authors |
|
|
* Copyright © 2016-2023 The Thingsboard Authors |
|
|
* |
|
|
* <p> |
|
|
* Licensed under the Apache License, Version 2.0 (the "License"); |
|
|
* Licensed under the Apache License, Version 2.0 (the "License"); |
|
|
* you may not use this file except in compliance with the License. |
|
|
* you may not use this file except in compliance with the License. |
|
|
* You may obtain a copy of the License at |
|
|
* You may obtain a copy of the License at |
|
|
* |
|
|
* <p> |
|
|
* http://www.apache.org/licenses/LICENSE-2.0
|
|
|
* http://www.apache.org/licenses/LICENSE-2.0
|
|
|
* |
|
|
* <p> |
|
|
* Unless required by applicable law or agreed to in writing, software |
|
|
* Unless required by applicable law or agreed to in writing, software |
|
|
* distributed under the License is distributed on an "AS IS" BASIS, |
|
|
* distributed under the License is distributed on an "AS IS" BASIS, |
|
|
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. |
|
|
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. |
|
|
@ -15,221 +15,459 @@ |
|
|
*/ |
|
|
*/ |
|
|
package org.thingsboard.rule.engine.metadata; |
|
|
package org.thingsboard.rule.engine.metadata; |
|
|
|
|
|
|
|
|
import com.google.common.collect.Lists; |
|
|
|
|
|
import com.google.common.util.concurrent.Futures; |
|
|
import com.google.common.util.concurrent.Futures; |
|
|
import org.junit.Before; |
|
|
import com.google.common.util.concurrent.ListenableFuture; |
|
|
import org.junit.Test; |
|
|
import org.jetbrains.annotations.NotNull; |
|
|
import org.junit.runner.RunWith; |
|
|
import org.junit.jupiter.api.BeforeEach; |
|
|
|
|
|
import org.junit.jupiter.api.Test; |
|
|
|
|
|
import org.junit.jupiter.api.extension.ExtendWith; |
|
|
import org.mockito.ArgumentCaptor; |
|
|
import org.mockito.ArgumentCaptor; |
|
|
import org.mockito.Mock; |
|
|
import org.mockito.Mock; |
|
|
import org.mockito.junit.MockitoJUnitRunner; |
|
|
import org.mockito.junit.jupiter.MockitoExtension; |
|
|
import org.thingsboard.common.util.JacksonUtil; |
|
|
import org.thingsboard.common.util.JacksonUtil; |
|
|
|
|
|
import org.thingsboard.common.util.ListeningExecutor; |
|
|
|
|
|
import org.thingsboard.rule.engine.api.TbContext; |
|
|
import org.thingsboard.rule.engine.api.TbNodeConfiguration; |
|
|
import org.thingsboard.rule.engine.api.TbNodeConfiguration; |
|
|
import org.thingsboard.rule.engine.api.TbNodeException; |
|
|
import org.thingsboard.rule.engine.api.TbNodeException; |
|
|
|
|
|
import org.thingsboard.rule.engine.data.RelationsQuery; |
|
|
|
|
|
import org.thingsboard.server.common.data.Customer; |
|
|
import org.thingsboard.server.common.data.Device; |
|
|
import org.thingsboard.server.common.data.Device; |
|
|
import org.thingsboard.server.common.data.User; |
|
|
import org.thingsboard.server.common.data.User; |
|
|
import org.thingsboard.server.common.data.asset.Asset; |
|
|
import org.thingsboard.server.common.data.asset.Asset; |
|
|
import org.thingsboard.server.common.data.id.AssetId; |
|
|
import org.thingsboard.server.common.data.id.AssetId; |
|
|
|
|
|
import org.thingsboard.server.common.data.id.CustomerId; |
|
|
import org.thingsboard.server.common.data.id.DeviceId; |
|
|
import org.thingsboard.server.common.data.id.DeviceId; |
|
|
import org.thingsboard.server.common.data.id.EntityId; |
|
|
import org.thingsboard.server.common.data.id.EntityId; |
|
|
|
|
|
import org.thingsboard.server.common.data.id.TenantId; |
|
|
import org.thingsboard.server.common.data.id.UserId; |
|
|
import org.thingsboard.server.common.data.id.UserId; |
|
|
import org.thingsboard.server.common.data.kv.AttributeKvEntry; |
|
|
import org.thingsboard.server.common.data.kv.AttributeKvEntry; |
|
|
import org.thingsboard.server.common.data.kv.BaseAttributeKvEntry; |
|
|
import org.thingsboard.server.common.data.kv.BaseAttributeKvEntry; |
|
|
|
|
|
import org.thingsboard.server.common.data.kv.BasicTsKvEntry; |
|
|
import org.thingsboard.server.common.data.kv.StringDataEntry; |
|
|
import org.thingsboard.server.common.data.kv.StringDataEntry; |
|
|
|
|
|
import org.thingsboard.server.common.data.kv.TsKvEntry; |
|
|
import org.thingsboard.server.common.data.relation.EntityRelation; |
|
|
import org.thingsboard.server.common.data.relation.EntityRelation; |
|
|
|
|
|
import org.thingsboard.server.common.data.relation.EntitySearchDirection; |
|
|
|
|
|
import org.thingsboard.server.common.data.relation.RelationEntityTypeFilter; |
|
|
import org.thingsboard.server.common.msg.TbMsg; |
|
|
import org.thingsboard.server.common.msg.TbMsg; |
|
|
import org.thingsboard.server.common.msg.TbMsgDataType; |
|
|
|
|
|
import org.thingsboard.server.common.msg.TbMsgMetaData; |
|
|
import org.thingsboard.server.common.msg.TbMsgMetaData; |
|
|
|
|
|
import org.thingsboard.server.dao.attributes.AttributesService; |
|
|
import org.thingsboard.server.dao.relation.RelationService; |
|
|
import org.thingsboard.server.dao.relation.RelationService; |
|
|
|
|
|
import org.thingsboard.server.dao.timeseries.TimeseriesService; |
|
|
|
|
|
|
|
|
import java.util.HashMap; |
|
|
import java.util.Collections; |
|
|
import java.util.List; |
|
|
import java.util.List; |
|
|
import java.util.Map; |
|
|
import java.util.Map; |
|
|
|
|
|
import java.util.NoSuchElementException; |
|
|
import java.util.UUID; |
|
|
import java.util.UUID; |
|
|
|
|
|
import java.util.concurrent.Callable; |
|
|
|
|
|
|
|
|
import static org.assertj.core.api.Assertions.assertThat; |
|
|
import static org.assertj.core.api.Assertions.assertThat; |
|
|
|
|
|
import static org.junit.jupiter.api.Assertions.assertEquals; |
|
|
|
|
|
import static org.junit.jupiter.api.Assertions.assertInstanceOf; |
|
|
import static org.junit.jupiter.api.Assertions.assertThrows; |
|
|
import static org.junit.jupiter.api.Assertions.assertThrows; |
|
|
import static org.mockito.ArgumentMatchers.any; |
|
|
import static org.mockito.ArgumentMatchers.any; |
|
|
import static org.mockito.ArgumentMatchers.anyCollection; |
|
|
import static org.mockito.ArgumentMatchers.anyList; |
|
|
import static org.mockito.ArgumentMatchers.eq; |
|
|
import static org.mockito.ArgumentMatchers.eq; |
|
|
|
|
|
import static org.mockito.Mockito.doReturn; |
|
|
|
|
|
import static org.mockito.Mockito.doThrow; |
|
|
import static org.mockito.Mockito.never; |
|
|
import static org.mockito.Mockito.never; |
|
|
import static org.mockito.Mockito.times; |
|
|
import static org.mockito.Mockito.times; |
|
|
import static org.mockito.Mockito.verify; |
|
|
import static org.mockito.Mockito.verify; |
|
|
import static org.mockito.Mockito.when; |
|
|
import static org.mockito.Mockito.when; |
|
|
import static org.thingsboard.server.common.data.DataConstants.SERVER_SCOPE; |
|
|
import static org.thingsboard.server.common.data.DataConstants.SERVER_SCOPE; |
|
|
|
|
|
|
|
|
@RunWith(MockitoJUnitRunner.class) |
|
|
@ExtendWith(MockitoExtension.class) |
|
|
public class TbGetRelatedAttributeNodeTest extends TbAbstractAttributeNodeTest { |
|
|
public class TbGetRelatedAttributeNodeTest { |
|
|
User user = new User(); |
|
|
|
|
|
Asset asset = new Asset(); |
|
|
private static final EntityId DUMMY_ENTITY_ID = new DeviceId(UUID.randomUUID()); |
|
|
Device device = new Device(); |
|
|
private static final TenantId TENANT_ID = new TenantId(UUID.randomUUID()); |
|
|
|
|
|
private static final ListeningExecutor DB_EXECUTOR = new ListeningExecutor() { |
|
|
|
|
|
@Override |
|
|
|
|
|
public <T> ListenableFuture<T> executeAsync(Callable<T> task) { |
|
|
|
|
|
try { |
|
|
|
|
|
return Futures.immediateFuture(task.call()); |
|
|
|
|
|
} catch (Exception e) { |
|
|
|
|
|
throw new RuntimeException(e); |
|
|
|
|
|
} |
|
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
|
|
@Override |
|
|
|
|
|
public void execute(@NotNull Runnable command) { |
|
|
|
|
|
command.run(); |
|
|
|
|
|
} |
|
|
|
|
|
}; |
|
|
@Mock |
|
|
@Mock |
|
|
private RelationService relationService; |
|
|
private TbContext ctxMock; |
|
|
|
|
|
@Mock |
|
|
|
|
|
private AttributesService attributesServiceMock; |
|
|
|
|
|
@Mock |
|
|
|
|
|
private TimeseriesService timeseriesServiceMock; |
|
|
|
|
|
@Mock |
|
|
|
|
|
private RelationService relationServiceMock; |
|
|
|
|
|
private TbGetRelatedAttributeNode node; |
|
|
|
|
|
private TbGetRelatedAttrNodeConfiguration config; |
|
|
|
|
|
private TbNodeConfiguration nodeConfiguration; |
|
|
private EntityRelation entityRelation; |
|
|
private EntityRelation entityRelation; |
|
|
|
|
|
private TbMsg msg; |
|
|
|
|
|
|
|
|
@Before |
|
|
@BeforeEach |
|
|
public void initDataForTests() throws TbNodeException { |
|
|
public void setUp() { |
|
|
init(new TbGetRelatedAttributeNode()); |
|
|
node = new TbGetRelatedAttributeNode(); |
|
|
|
|
|
config = new TbGetRelatedAttrNodeConfiguration().defaultConfiguration(); |
|
|
|
|
|
nodeConfiguration = new TbNodeConfiguration(JacksonUtil.valueToTree(config)); |
|
|
entityRelation = new EntityRelation(); |
|
|
entityRelation = new EntityRelation(); |
|
|
entityRelation.setTo(customerId); |
|
|
|
|
|
entityRelation.setType(EntityRelation.CONTAINS_TYPE); |
|
|
|
|
|
when(ctx.getRelationService()).thenReturn(relationService); |
|
|
|
|
|
|
|
|
|
|
|
user.setCustomerId(customerId); |
|
|
|
|
|
user.setId(new UserId(UUID.randomUUID())); |
|
|
|
|
|
entityRelation.setFrom(user.getId()); |
|
|
|
|
|
|
|
|
|
|
|
asset.setCustomerId(customerId); |
|
|
|
|
|
asset.setId(new AssetId(UUID.randomUUID())); |
|
|
|
|
|
|
|
|
|
|
|
device.setCustomerId(customerId); |
|
|
|
|
|
device.setId(new DeviceId(UUID.randomUUID())); |
|
|
|
|
|
} |
|
|
} |
|
|
|
|
|
|
|
|
@Override |
|
|
@Test |
|
|
protected TbAbstractGetEntityAttrNode getEmptyNode() { |
|
|
public void givenConfigWithNullFetchTo_whenInit_thenException() { |
|
|
return new TbGetRelatedAttributeNode(); |
|
|
// GIVEN
|
|
|
} |
|
|
config.setFetchTo(null); |
|
|
|
|
|
nodeConfiguration = new TbNodeConfiguration(JacksonUtil.valueToTree(config)); |
|
|
|
|
|
|
|
|
@Override |
|
|
// WHEN
|
|
|
TbGetEntityAttrNodeConfiguration getTbNodeConfig() { |
|
|
var exception = assertThrows(TbNodeException.class, () -> node.init(ctxMock, nodeConfiguration)); |
|
|
return getConfig(false); |
|
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
|
|
@Override |
|
|
// THEN
|
|
|
TbGetEntityAttrNodeConfiguration getTbNodeConfigForTelemetry() { |
|
|
assertThat(exception.getMessage()).isEqualTo("FetchTo cannot be null!"); |
|
|
return getConfig(true); |
|
|
verify(ctxMock, never()).tellSuccess(any()); |
|
|
} |
|
|
} |
|
|
|
|
|
|
|
|
private TbGetEntityAttrNodeConfiguration getConfig(boolean isTelemetry) { |
|
|
@Test |
|
|
TbGetRelatedAttrNodeConfiguration config = new TbGetRelatedAttrNodeConfiguration(); |
|
|
public void givenDefaultConfig_whenInit_thenOK() throws TbNodeException { |
|
|
config = config.defaultConfiguration(); |
|
|
// GIVEN
|
|
|
Map<String, String> conf = new HashMap<>(); |
|
|
|
|
|
conf.put(keyAttrConf, valueAttrConf); |
|
|
// WHEN
|
|
|
config.setAttrMapping(conf); |
|
|
node.init(ctxMock, nodeConfiguration); |
|
|
config.setTelemetry(isTelemetry); |
|
|
|
|
|
config.setFetchTo(FetchTo.METADATA); |
|
|
// THEN
|
|
|
return config; |
|
|
var nodeConfig = (TbGetRelatedAttrNodeConfiguration) node.config; |
|
|
|
|
|
assertThat(nodeConfig).isEqualTo(config); |
|
|
|
|
|
assertThat(nodeConfig.getAttrMapping()).isEqualTo(Map.of("serialNumber", "sn")); |
|
|
|
|
|
assertThat(nodeConfig.isTelemetry()).isEqualTo(false); |
|
|
|
|
|
assertThat(node.fetchTo).isEqualTo(FetchTo.METADATA); |
|
|
|
|
|
|
|
|
|
|
|
var relationsQuery = new RelationsQuery(); |
|
|
|
|
|
var relationEntityTypeFilter = new RelationEntityTypeFilter(EntityRelation.CONTAINS_TYPE, Collections.emptyList()); |
|
|
|
|
|
relationsQuery.setDirection(EntitySearchDirection.FROM); |
|
|
|
|
|
relationsQuery.setMaxLevel(1); |
|
|
|
|
|
relationsQuery.setFilters(Collections.singletonList(relationEntityTypeFilter)); |
|
|
|
|
|
|
|
|
|
|
|
assertThat(nodeConfig.getRelationsQuery()).isEqualTo(relationsQuery); |
|
|
} |
|
|
} |
|
|
|
|
|
|
|
|
@Override |
|
|
@Test |
|
|
EntityId getEntityId() { |
|
|
public void givenCustomConfig_whenInit_thenOK() throws TbNodeException { |
|
|
return customerId; |
|
|
// GIVEN
|
|
|
|
|
|
config.setAttrMapping(Map.of( |
|
|
|
|
|
"sourceAttr1", "targetKey1", |
|
|
|
|
|
"sourceAttr2", "targetKey2", |
|
|
|
|
|
"sourceAttr3", "targetKey3")); |
|
|
|
|
|
config.setTelemetry(true); |
|
|
|
|
|
config.setFetchTo(FetchTo.DATA); |
|
|
|
|
|
|
|
|
|
|
|
var relationsQuery = new RelationsQuery(); |
|
|
|
|
|
var relationEntityTypeFilter = new RelationEntityTypeFilter(EntityRelation.CONTAINS_TYPE, Collections.emptyList()); |
|
|
|
|
|
relationsQuery.setDirection(EntitySearchDirection.FROM); |
|
|
|
|
|
relationsQuery.setMaxLevel(1); |
|
|
|
|
|
relationsQuery.setFilters(Collections.singletonList(relationEntityTypeFilter)); |
|
|
|
|
|
|
|
|
|
|
|
config.setRelationsQuery(relationsQuery); |
|
|
|
|
|
nodeConfiguration = new TbNodeConfiguration(JacksonUtil.valueToTree(config)); |
|
|
|
|
|
|
|
|
|
|
|
// WHEN
|
|
|
|
|
|
node.init(ctxMock, nodeConfiguration); |
|
|
|
|
|
|
|
|
|
|
|
// THEN
|
|
|
|
|
|
var nodeConfig = (TbGetRelatedAttrNodeConfiguration) node.config; |
|
|
|
|
|
assertThat(nodeConfig).isEqualTo(config); |
|
|
|
|
|
assertThat(nodeConfig.getAttrMapping()).isEqualTo(Map.of( |
|
|
|
|
|
"sourceAttr1", "targetKey1", |
|
|
|
|
|
"sourceAttr2", "targetKey2", |
|
|
|
|
|
"sourceAttr3", "targetKey3" |
|
|
|
|
|
)); |
|
|
|
|
|
assertThat(nodeConfig.isTelemetry()).isEqualTo(true); |
|
|
|
|
|
assertThat(node.fetchTo).isEqualTo(FetchTo.DATA); |
|
|
|
|
|
assertThat(nodeConfig.getRelationsQuery()).isEqualTo(relationsQuery); |
|
|
} |
|
|
} |
|
|
|
|
|
|
|
|
@Test |
|
|
@Test |
|
|
public void errorThrownIfFetchToIsNull() { |
|
|
public void givenEmptyAttributesMapping_whenInit_thenException() { |
|
|
var node = new TbGetRelatedAttributeNode(); |
|
|
// SETUP
|
|
|
var config = new TbGetRelatedAttrNodeConfiguration().defaultConfiguration(); |
|
|
var expectedExceptionMessage = "At least one attribute mapping should be specified!"; |
|
|
config.setFetchTo(null); |
|
|
|
|
|
var nodeConfiguration = new TbNodeConfiguration(JacksonUtil.valueToTree(config)); |
|
|
// GIVEN
|
|
|
|
|
|
config.setAttrMapping(Collections.emptyMap()); |
|
|
|
|
|
nodeConfiguration = new TbNodeConfiguration(JacksonUtil.valueToTree(config)); |
|
|
|
|
|
|
|
|
var exception = assertThrows(TbNodeException.class, () -> node.init(ctx, nodeConfiguration)); |
|
|
// WHEN
|
|
|
|
|
|
var exception = assertThrows(TbNodeException.class, () -> node.init(ctxMock, nodeConfiguration)); |
|
|
|
|
|
|
|
|
assertThat(exception.getMessage()).isEqualTo("FetchTo cannot be NULL!"); |
|
|
// THEN
|
|
|
verify(ctx, never()).tellSuccess(any()); |
|
|
assertThat(exception.getMessage()).isEqualTo(expectedExceptionMessage); |
|
|
|
|
|
verify(ctxMock, never()).tellSuccess(any()); |
|
|
} |
|
|
} |
|
|
|
|
|
|
|
|
@Test |
|
|
@Test |
|
|
public void errorThrownIfMsgDataIsNotAnObjectAndFetchToData() { |
|
|
public void givenMsgDataIsNotAnJsonObjectAndFetchToData_whenOnMsg_thenException() { |
|
|
|
|
|
// GIVEN
|
|
|
node.fetchTo = FetchTo.DATA; |
|
|
node.fetchTo = FetchTo.DATA; |
|
|
node.config.setFetchTo(FetchTo.DATA); |
|
|
msg = TbMsg.newMsg("POST_TELEMETRY_REQUEST", DUMMY_ENTITY_ID, new TbMsgMetaData(), "[]"); |
|
|
msg = TbMsg.newMsg("SOME_MESSAGE_TYPE", new DeviceId(UUID.randomUUID()), new TbMsgMetaData(), "[]"); |
|
|
|
|
|
|
|
|
|
|
|
var exception = assertThrows(IllegalArgumentException.class, () -> node.onMsg(ctx, msg)); |
|
|
// WHEN
|
|
|
|
|
|
var exception = assertThrows(IllegalArgumentException.class, () -> node.onMsg(ctxMock, msg)); |
|
|
|
|
|
|
|
|
|
|
|
// THEN
|
|
|
assertThat(exception.getMessage()).isEqualTo("Message body is not an object!"); |
|
|
assertThat(exception.getMessage()).isEqualTo("Message body is not an object!"); |
|
|
verify(ctx, never()).tellSuccess(any()); |
|
|
verify(ctxMock, never()).tellSuccess(any()); |
|
|
} |
|
|
} |
|
|
|
|
|
|
|
|
@Test |
|
|
@Test |
|
|
public void errorThrownIfCannotLoadAttributes() { |
|
|
public void givenEntityThatDoesNotBelongToTheCurrentTenant_whenOnMsg_thenException() { |
|
|
entityRelation.setFrom(user.getId()); |
|
|
// SETUP
|
|
|
entityRelation.setTo(customerId); |
|
|
var expectedExceptionMessage = "Entity with id: '" + DUMMY_ENTITY_ID + |
|
|
when(relationService.findByQuery(any(), any())).thenReturn(Futures.immediateFuture(List.of(entityRelation))); |
|
|
"' specified in the configuration doesn't belong to the current tenant."; |
|
|
errorThrownIfCannotLoadAttributes(user); |
|
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
|
|
@Test |
|
|
// GIVEN
|
|
|
public void errorThrownIfCannotLoadAttributesAsync() { |
|
|
doThrow(new RuntimeException(expectedExceptionMessage)).when(ctxMock).checkTenantEntity(DUMMY_ENTITY_ID); |
|
|
entityRelation.setFrom(user.getId()); |
|
|
msg = TbMsg.newMsg("POST_TELEMETRY_REQUEST", DUMMY_ENTITY_ID, new TbMsgMetaData(), "{}"); |
|
|
entityRelation.setTo(customerId); |
|
|
|
|
|
when(relationService.findByQuery(any(), any())).thenReturn(Futures.immediateFuture(List.of(entityRelation))); |
|
|
|
|
|
errorThrownIfCannotLoadAttributesAsync(user); |
|
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
|
|
@Test |
|
|
// WHEN
|
|
|
public void failedChainUsedIfCustomerCannotBeFound() { |
|
|
var exception = assertThrows(RuntimeException.class, () -> node.onMsg(ctxMock, msg)); |
|
|
entityRelation.setFrom(customerId); |
|
|
|
|
|
entityRelation.setTo(null); |
|
|
// THEN
|
|
|
when(relationService.findByQuery(any(), any())).thenReturn(Futures.immediateFuture(List.of(entityRelation))); |
|
|
assertThat(exception.getMessage()).isEqualTo(expectedExceptionMessage); |
|
|
failedChainUsedIfCustomerCannotBeFound(user); |
|
|
verify(ctxMock, never()).tellSuccess(any()); |
|
|
} |
|
|
} |
|
|
|
|
|
|
|
|
@Test |
|
|
@Test |
|
|
public void customerAttributeAddedInMetadata() { |
|
|
public void givenDidNotFindEntity_whenOnMsg_thenShouldTellFailure() { |
|
|
entityRelation.setFrom(customerId); |
|
|
// GIVEN
|
|
|
entityRelation.setTo(customerId); |
|
|
prepareMsgAndConfig(FetchTo.METADATA, false, DUMMY_ENTITY_ID); |
|
|
when(relationService.findByQuery(any(), any())).thenReturn(Futures.immediateFuture(List.of(entityRelation))); |
|
|
|
|
|
entityAttributeAddedInMetadata(customerId, "CUSTOMER"); |
|
|
when(ctxMock.getDbCallbackExecutor()).thenReturn(DB_EXECUTOR); |
|
|
|
|
|
when(ctxMock.getRelationService()).thenReturn(relationServiceMock); |
|
|
|
|
|
doReturn(Futures.immediateFuture(null)).when(relationServiceMock).findByQuery(any(), any()); |
|
|
|
|
|
|
|
|
|
|
|
// WHEN
|
|
|
|
|
|
node.onMsg(ctxMock, msg); |
|
|
|
|
|
|
|
|
|
|
|
// THEN
|
|
|
|
|
|
var actualMessageCaptor = ArgumentCaptor.forClass(TbMsg.class); |
|
|
|
|
|
var actualExceptionCaptor = ArgumentCaptor.forClass(Throwable.class); |
|
|
|
|
|
|
|
|
|
|
|
verify(ctxMock, never()).tellSuccess(any()); |
|
|
|
|
|
verify(ctxMock, times(1)) |
|
|
|
|
|
.tellFailure(actualMessageCaptor.capture(), actualExceptionCaptor.capture()); |
|
|
|
|
|
|
|
|
|
|
|
var actualMessage = actualMessageCaptor.getValue(); |
|
|
|
|
|
var actualException = actualExceptionCaptor.getValue(); |
|
|
|
|
|
|
|
|
|
|
|
var expectedExceptionMessage = "Failed to find related entity to message originator using relation query specified in the configuration!"; |
|
|
|
|
|
|
|
|
|
|
|
assertEquals(msg, actualMessage); |
|
|
|
|
|
assertEquals(expectedExceptionMessage, actualException.getMessage()); |
|
|
|
|
|
assertInstanceOf(NoSuchElementException.class, actualException); |
|
|
} |
|
|
} |
|
|
|
|
|
|
|
|
@Test |
|
|
@Test |
|
|
public void customerAttributeAddedInData() { |
|
|
public void givenFetchAttributesToData_whenOnMsg_thenShouldFetchAttributesToData() { |
|
|
node.fetchTo = FetchTo.DATA; |
|
|
// GIVEN
|
|
|
node.config.setFetchTo(FetchTo.DATA); |
|
|
var customer = new Customer(new CustomerId(UUID.randomUUID())); |
|
|
|
|
|
var user = new User(new UserId(UUID.randomUUID())); |
|
|
|
|
|
|
|
|
|
|
|
prepareMsgAndConfig(FetchTo.DATA, false, user.getId()); |
|
|
|
|
|
|
|
|
entityRelation.setFrom(customerId); |
|
|
entityRelation.setFrom(user.getId()); |
|
|
entityRelation.setTo(customerId); |
|
|
entityRelation.setTo(customer.getId()); |
|
|
when(relationService.findByQuery(any(), any())).thenReturn(Futures.immediateFuture(List.of(entityRelation))); |
|
|
entityRelation.setType(EntityRelation.CONTAINS_TYPE); |
|
|
|
|
|
|
|
|
msg = TbMsg.newMsg("CUSTOMER", customerId, new TbMsgMetaData(metaData), TbMsgDataType.JSON, "{}", ruleChainId, ruleNodeId); |
|
|
when(ctxMock.getRelationService()).thenReturn(relationServiceMock); |
|
|
|
|
|
doReturn(Futures.immediateFuture(List.of(entityRelation))).when(relationServiceMock).findByQuery(eq(TENANT_ID), any()); |
|
|
|
|
|
|
|
|
List<AttributeKvEntry> attributes = Lists.newArrayList(new BaseAttributeKvEntry(new StringDataEntry("temperature", "high"), 1L)); |
|
|
when(ctxMock.getTenantId()).thenReturn(TENANT_ID); |
|
|
|
|
|
when(ctxMock.getAttributesService()).thenReturn(attributesServiceMock); |
|
|
|
|
|
|
|
|
when(ctx.getAttributesService()).thenReturn(attributesService); |
|
|
List<AttributeKvEntry> attributes = List.of( |
|
|
when(attributesService.find(any(), eq(customerId), eq(SERVER_SCOPE), anyCollection())) |
|
|
new BaseAttributeKvEntry(new StringDataEntry("sourceKey1", "sourceValue1"), 1L), |
|
|
|
|
|
new BaseAttributeKvEntry(new StringDataEntry("sourceKey2", "sourceValue2"), 2L), |
|
|
|
|
|
new BaseAttributeKvEntry(new StringDataEntry("sourceKey3", "sourceValue3"), 3L) |
|
|
|
|
|
); |
|
|
|
|
|
when(attributesServiceMock.find(eq(TENANT_ID), eq(customer.getId()), eq(SERVER_SCOPE), anyList())) |
|
|
.thenReturn(Futures.immediateFuture(attributes)); |
|
|
.thenReturn(Futures.immediateFuture(attributes)); |
|
|
|
|
|
when(ctxMock.getDbCallbackExecutor()).thenReturn(DB_EXECUTOR); |
|
|
|
|
|
|
|
|
node.onMsg(ctx, msg); |
|
|
// WHEN
|
|
|
|
|
|
node.onMsg(ctxMock, msg); |
|
|
|
|
|
|
|
|
|
|
|
// THEN
|
|
|
var actualMessageCaptor = ArgumentCaptor.forClass(TbMsg.class); |
|
|
var actualMessageCaptor = ArgumentCaptor.forClass(TbMsg.class); |
|
|
verify(ctx, times(1)).tellSuccess(actualMessageCaptor.capture()); |
|
|
verify(ctxMock, times(1)).tellSuccess(actualMessageCaptor.capture()); |
|
|
|
|
|
verify(ctxMock, never()).tellFailure(any(), any()); |
|
|
|
|
|
|
|
|
var expectedMsgData = "{\"answer\":\"high\"}"; |
|
|
var expectedMsgData = "{\"temp\":42," + |
|
|
|
|
|
"\"humidity\":77," + |
|
|
|
|
|
"\"messageBodyPattern1\":\"targetKey2\"," + |
|
|
|
|
|
"\"messageBodyPattern2\":\"sourceKey3\"," + |
|
|
|
|
|
"\"targetKey1\":\"sourceValue1\"," + |
|
|
|
|
|
"\"targetKey2\":\"sourceValue2\"," + |
|
|
|
|
|
"\"targetKey3\":\"sourceValue3\"}"; |
|
|
|
|
|
|
|
|
assertThat(actualMessageCaptor.getValue().getData()).isEqualTo(expectedMsgData); |
|
|
assertThat(actualMessageCaptor.getValue().getData()).isEqualTo(expectedMsgData); |
|
|
|
|
|
assertThat(actualMessageCaptor.getValue().getMetaData()).isEqualTo(msg.getMetaData()); |
|
|
} |
|
|
} |
|
|
|
|
|
|
|
|
@Test |
|
|
@Test |
|
|
public void usersCustomerAttributesFetched() { |
|
|
public void givenFetchAttributesToMetaData_whenOnMsg_thenShouldFetchAttributesToMetaData() { |
|
|
entityRelation.setFrom(user.getId()); |
|
|
// GIVEN
|
|
|
entityRelation.setTo(customerId); |
|
|
var customer = new Customer(new CustomerId(UUID.randomUUID())); |
|
|
when(relationService.findByQuery(any(), any())).thenReturn(Futures.immediateFuture(List.of(entityRelation))); |
|
|
|
|
|
usersCustomerAttributesFetched(user); |
|
|
prepareMsgAndConfig(FetchTo.METADATA, false, customer.getId()); |
|
|
|
|
|
|
|
|
|
|
|
entityRelation.setFrom(customer.getId()); |
|
|
|
|
|
entityRelation.setTo(customer.getId()); |
|
|
|
|
|
entityRelation.setType(EntityRelation.CONTAINS_TYPE); |
|
|
|
|
|
|
|
|
|
|
|
when(ctxMock.getRelationService()).thenReturn(relationServiceMock); |
|
|
|
|
|
doReturn(Futures.immediateFuture(List.of(entityRelation))).when(relationServiceMock).findByQuery(eq(TENANT_ID), any()); |
|
|
|
|
|
|
|
|
|
|
|
when(ctxMock.getTenantId()).thenReturn(TENANT_ID); |
|
|
|
|
|
when(ctxMock.getAttributesService()).thenReturn(attributesServiceMock); |
|
|
|
|
|
List<AttributeKvEntry> attributes = List.of( |
|
|
|
|
|
new BaseAttributeKvEntry(new StringDataEntry("sourceKey1", "sourceValue1"), 1L), |
|
|
|
|
|
new BaseAttributeKvEntry(new StringDataEntry("sourceKey2", "sourceValue2"), 2L), |
|
|
|
|
|
new BaseAttributeKvEntry(new StringDataEntry("sourceKey3", "sourceValue3"), 3L) |
|
|
|
|
|
); |
|
|
|
|
|
when(attributesServiceMock.find(eq(TENANT_ID), eq(customer.getId()), eq(SERVER_SCOPE), anyList())) |
|
|
|
|
|
.thenReturn(Futures.immediateFuture(attributes)); |
|
|
|
|
|
when(ctxMock.getDbCallbackExecutor()).thenReturn(DB_EXECUTOR); |
|
|
|
|
|
|
|
|
|
|
|
// WHEN
|
|
|
|
|
|
node.onMsg(ctxMock, msg); |
|
|
|
|
|
|
|
|
|
|
|
// THEN
|
|
|
|
|
|
var actualMessageCaptor = ArgumentCaptor.forClass(TbMsg.class); |
|
|
|
|
|
verify(ctxMock, times(1)).tellSuccess(actualMessageCaptor.capture()); |
|
|
|
|
|
verify(ctxMock, never()).tellFailure(any(), any()); |
|
|
|
|
|
|
|
|
|
|
|
var expectedMsgMetaData = new TbMsgMetaData(Map.of( |
|
|
|
|
|
"metaDataPattern1", "sourceKey2", |
|
|
|
|
|
"metaDataPattern2", "targetKey3", |
|
|
|
|
|
"targetKey1", "sourceValue1", |
|
|
|
|
|
"targetKey2", "sourceValue2", |
|
|
|
|
|
"targetKey3", "sourceValue3" |
|
|
|
|
|
)); |
|
|
|
|
|
|
|
|
|
|
|
assertThat(actualMessageCaptor.getValue().getData()).isEqualTo(msg.getData()); |
|
|
|
|
|
assertThat(actualMessageCaptor.getValue().getMetaData()).isEqualTo(expectedMsgMetaData); |
|
|
} |
|
|
} |
|
|
|
|
|
|
|
|
@Test |
|
|
@Test |
|
|
public void assetsCustomerAttributesFetched() { |
|
|
public void givenFetchTelemetryToData_whenOnMsg_thenShouldFetchTelemetryToData() { |
|
|
|
|
|
// GIVEN
|
|
|
|
|
|
var customer = new Customer(new CustomerId(UUID.randomUUID())); |
|
|
|
|
|
var asset = new Asset(new AssetId(UUID.randomUUID())); |
|
|
|
|
|
|
|
|
|
|
|
prepareMsgAndConfig(FetchTo.DATA, true, asset.getId()); |
|
|
|
|
|
|
|
|
entityRelation.setFrom(asset.getId()); |
|
|
entityRelation.setFrom(asset.getId()); |
|
|
entityRelation.setTo(customerId); |
|
|
entityRelation.setTo(customer.getId()); |
|
|
when(relationService.findByQuery(any(), any())).thenReturn(Futures.immediateFuture(List.of(entityRelation))); |
|
|
entityRelation.setType(EntityRelation.CONTAINS_TYPE); |
|
|
assetsCustomerAttributesFetched(asset); |
|
|
|
|
|
|
|
|
when(ctxMock.getRelationService()).thenReturn(relationServiceMock); |
|
|
|
|
|
doReturn(Futures.immediateFuture(List.of(entityRelation))).when(relationServiceMock).findByQuery(eq(TENANT_ID), any()); |
|
|
|
|
|
|
|
|
|
|
|
when(ctxMock.getTenantId()).thenReturn(TENANT_ID); |
|
|
|
|
|
when(ctxMock.getTimeseriesService()).thenReturn(timeseriesServiceMock); |
|
|
|
|
|
List<TsKvEntry> timeseries = List.of( |
|
|
|
|
|
new BasicTsKvEntry(1L, new StringDataEntry("sourceKey1", "sourceValue1")), |
|
|
|
|
|
new BasicTsKvEntry(1L, new StringDataEntry("sourceKey2", "sourceValue2")), |
|
|
|
|
|
new BasicTsKvEntry(1L, new StringDataEntry("sourceKey3", "sourceValue3")) |
|
|
|
|
|
); |
|
|
|
|
|
when(timeseriesServiceMock.findLatest(eq(TENANT_ID), eq(customer.getId()), anyList())) |
|
|
|
|
|
.thenReturn(Futures.immediateFuture(timeseries)); |
|
|
|
|
|
when(ctxMock.getDbCallbackExecutor()).thenReturn(DB_EXECUTOR); |
|
|
|
|
|
|
|
|
|
|
|
// WHEN
|
|
|
|
|
|
node.onMsg(ctxMock, msg); |
|
|
|
|
|
|
|
|
|
|
|
// THEN
|
|
|
|
|
|
var actualMessageCaptor = ArgumentCaptor.forClass(TbMsg.class); |
|
|
|
|
|
verify(ctxMock, times(1)).tellSuccess(actualMessageCaptor.capture()); |
|
|
|
|
|
verify(ctxMock, never()).tellFailure(any(), any()); |
|
|
|
|
|
|
|
|
|
|
|
var expectedMsgData = "{\"temp\":42," + |
|
|
|
|
|
"\"humidity\":77," + |
|
|
|
|
|
"\"messageBodyPattern1\":\"targetKey2\"," + |
|
|
|
|
|
"\"messageBodyPattern2\":\"sourceKey3\"," + |
|
|
|
|
|
"\"targetKey1\":\"sourceValue1\"," + |
|
|
|
|
|
"\"targetKey2\":\"sourceValue2\"," + |
|
|
|
|
|
"\"targetKey3\":\"sourceValue3\"}"; |
|
|
|
|
|
|
|
|
|
|
|
assertThat(actualMessageCaptor.getValue().getData()).isEqualTo(expectedMsgData); |
|
|
|
|
|
assertThat(actualMessageCaptor.getValue().getMetaData()).isEqualTo(msg.getMetaData()); |
|
|
} |
|
|
} |
|
|
|
|
|
|
|
|
@Test |
|
|
@Test |
|
|
public void deviceCustomerAttributesFetched() { |
|
|
public void givenFetchTelemetryToMetaData_whenOnMsg_thenShouldFetchTelemetryToMetaData() { |
|
|
|
|
|
// GIVEN
|
|
|
|
|
|
var customer = new Customer(new CustomerId(UUID.randomUUID())); |
|
|
|
|
|
var device = new Device(new DeviceId(UUID.randomUUID())); |
|
|
|
|
|
|
|
|
|
|
|
prepareMsgAndConfig(FetchTo.METADATA, true, device.getId()); |
|
|
|
|
|
|
|
|
entityRelation.setFrom(device.getId()); |
|
|
entityRelation.setFrom(device.getId()); |
|
|
entityRelation.setTo(customerId); |
|
|
entityRelation.setTo(customer.getId()); |
|
|
when(relationService.findByQuery(any(), any())).thenReturn(Futures.immediateFuture(List.of(entityRelation))); |
|
|
entityRelation.setType(EntityRelation.CONTAINS_TYPE); |
|
|
deviceCustomerAttributesFetched(device); |
|
|
|
|
|
|
|
|
when(ctxMock.getRelationService()).thenReturn(relationServiceMock); |
|
|
|
|
|
doReturn(Futures.immediateFuture(List.of(entityRelation))).when(relationServiceMock).findByQuery(eq(TENANT_ID), any()); |
|
|
|
|
|
|
|
|
|
|
|
when(ctxMock.getTenantId()).thenReturn(TENANT_ID); |
|
|
|
|
|
when(ctxMock.getTimeseriesService()).thenReturn(timeseriesServiceMock); |
|
|
|
|
|
List<TsKvEntry> timeseries = List.of( |
|
|
|
|
|
new BasicTsKvEntry(1L, new StringDataEntry("sourceKey1", "sourceValue1")), |
|
|
|
|
|
new BasicTsKvEntry(1L, new StringDataEntry("sourceKey2", "sourceValue2")), |
|
|
|
|
|
new BasicTsKvEntry(1L, new StringDataEntry("sourceKey3", "sourceValue3")) |
|
|
|
|
|
); |
|
|
|
|
|
when(timeseriesServiceMock.findLatest(eq(TENANT_ID), eq(customer.getId()), anyList())) |
|
|
|
|
|
.thenReturn(Futures.immediateFuture(timeseries)); |
|
|
|
|
|
when(ctxMock.getDbCallbackExecutor()).thenReturn(DB_EXECUTOR); |
|
|
|
|
|
|
|
|
|
|
|
// WHEN
|
|
|
|
|
|
node.onMsg(ctxMock, msg); |
|
|
|
|
|
|
|
|
|
|
|
// THEN
|
|
|
|
|
|
var actualMessageCaptor = ArgumentCaptor.forClass(TbMsg.class); |
|
|
|
|
|
verify(ctxMock, times(1)).tellSuccess(actualMessageCaptor.capture()); |
|
|
|
|
|
verify(ctxMock, never()).tellFailure(any(), any()); |
|
|
|
|
|
|
|
|
|
|
|
var expectedMsgMetaData = new TbMsgMetaData(Map.of( |
|
|
|
|
|
"metaDataPattern1", "sourceKey2", |
|
|
|
|
|
"metaDataPattern2", "targetKey3", |
|
|
|
|
|
"targetKey1", "sourceValue1", |
|
|
|
|
|
"targetKey2", "sourceValue2", |
|
|
|
|
|
"targetKey3", "sourceValue3" |
|
|
|
|
|
)); |
|
|
|
|
|
|
|
|
|
|
|
assertThat(actualMessageCaptor.getValue().getData()).isEqualTo(msg.getData()); |
|
|
|
|
|
assertThat(actualMessageCaptor.getValue().getMetaData()).isEqualTo(expectedMsgMetaData); |
|
|
} |
|
|
} |
|
|
|
|
|
|
|
|
@Test |
|
|
private void prepareMsgAndConfig(FetchTo fetchTo, boolean isTelemetry, EntityId entityId) { |
|
|
public void deviceCustomerTelemetryFetched() throws TbNodeException { |
|
|
config.setAttrMapping(Map.of( |
|
|
entityRelation.setFrom(device.getId()); |
|
|
"sourceKey1", "targetKey1", |
|
|
entityRelation.setTo(customerId); |
|
|
"${metaDataPattern1}", "$[messageBodyPattern1]", |
|
|
when(relationService.findByQuery(any(), any())).thenReturn(Futures.immediateFuture(List.of(entityRelation))); |
|
|
"$[messageBodyPattern2]", "${metaDataPattern2}")); |
|
|
deviceCustomerTelemetryFetched(device); |
|
|
config.setTelemetry(isTelemetry); |
|
|
|
|
|
config.setFetchTo(fetchTo); |
|
|
|
|
|
|
|
|
|
|
|
node.config = config; |
|
|
|
|
|
node.fetchTo = fetchTo; |
|
|
|
|
|
var msgMetaData = new TbMsgMetaData(); |
|
|
|
|
|
msgMetaData.putValue("metaDataPattern1", "sourceKey2"); |
|
|
|
|
|
msgMetaData.putValue("metaDataPattern2", "targetKey3"); |
|
|
|
|
|
var msgData = "{\"temp\":42,\"humidity\":77,\"messageBodyPattern1\":\"targetKey2\",\"messageBodyPattern2\":\"sourceKey3\"}"; |
|
|
|
|
|
msg = TbMsg.newMsg("POST_TELEMETRY_REQUEST", entityId, msgMetaData, msgData); |
|
|
} |
|
|
} |
|
|
|
|
|
|
|
|
} |
|
|
} |
|
|
|