committed by
GitHub
38 changed files with 3252 additions and 1005 deletions
@ -0,0 +1,452 @@ |
|||
/** |
|||
* Copyright © 2016-2023 The Thingsboard Authors |
|||
* |
|||
* Licensed under the Apache License, Version 2.0 (the "License"); |
|||
* you may not use this file except in compliance with the License. |
|||
* You may obtain a copy of the License at |
|||
* |
|||
* http://www.apache.org/licenses/LICENSE-2.0
|
|||
* |
|||
* Unless required by applicable law or agreed to in writing, software |
|||
* distributed under the License is distributed on an "AS IS" BASIS, |
|||
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. |
|||
* See the License for the specific language governing permissions and |
|||
* limitations under the License. |
|||
*/ |
|||
package org.thingsboard.rule.engine.metadata; |
|||
|
|||
import com.google.common.util.concurrent.Futures; |
|||
import com.google.common.util.concurrent.ListenableFuture; |
|||
import lombok.RequiredArgsConstructor; |
|||
import org.jetbrains.annotations.NotNull; |
|||
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.ArgumentMatcher; |
|||
import org.mockito.Mock; |
|||
import org.mockito.junit.jupiter.MockitoExtension; |
|||
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.TbNodeException; |
|||
import org.thingsboard.server.common.data.id.DeviceId; |
|||
import org.thingsboard.server.common.data.id.TenantId; |
|||
import org.thingsboard.server.common.data.kv.BasicTsKvEntry; |
|||
import org.thingsboard.server.common.data.kv.BooleanDataEntry; |
|||
import org.thingsboard.server.common.data.kv.DoubleDataEntry; |
|||
import org.thingsboard.server.common.data.kv.JsonDataEntry; |
|||
import org.thingsboard.server.common.data.kv.LongDataEntry; |
|||
import org.thingsboard.server.common.data.kv.StringDataEntry; |
|||
import org.thingsboard.server.common.data.kv.TsKvEntry; |
|||
import org.thingsboard.server.common.msg.TbMsg; |
|||
import org.thingsboard.server.common.msg.TbMsgMetaData; |
|||
import org.thingsboard.server.dao.timeseries.TimeseriesService; |
|||
|
|||
import java.util.List; |
|||
import java.util.UUID; |
|||
import java.util.concurrent.Callable; |
|||
|
|||
import static org.junit.jupiter.api.Assertions.assertEquals; |
|||
import static org.junit.jupiter.api.Assertions.assertFalse; |
|||
import static org.junit.jupiter.api.Assertions.assertInstanceOf; |
|||
import static org.junit.jupiter.api.Assertions.assertTrue; |
|||
import static org.mockito.ArgumentMatchers.any; |
|||
import static org.mockito.ArgumentMatchers.anyList; |
|||
import static org.mockito.ArgumentMatchers.anySet; |
|||
import static org.mockito.ArgumentMatchers.anyString; |
|||
import static org.mockito.ArgumentMatchers.argThat; |
|||
import static org.mockito.ArgumentMatchers.eq; |
|||
import static org.mockito.Mockito.never; |
|||
import static org.mockito.Mockito.reset; |
|||
import static org.mockito.Mockito.times; |
|||
import static org.mockito.Mockito.verify; |
|||
import static org.mockito.Mockito.when; |
|||
|
|||
@ExtendWith(MockitoExtension.class) |
|||
public class CalculateDeltaNodeTest { |
|||
|
|||
private static final DeviceId DUMMY_DEVICE_ORIGINATOR = new DeviceId(UUID.randomUUID()); |
|||
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 |
|||
private TbContext ctxMock; |
|||
@Mock |
|||
private TimeseriesService timeseriesServiceMock; |
|||
private CalculateDeltaNode node; |
|||
private CalculateDeltaNodeConfiguration config; |
|||
private TbNodeConfiguration nodeConfiguration; |
|||
|
|||
@BeforeEach |
|||
public void setUp() throws TbNodeException { |
|||
node = new CalculateDeltaNode(); |
|||
config = new CalculateDeltaNodeConfiguration().defaultConfiguration(); |
|||
nodeConfiguration = new TbNodeConfiguration(JacksonUtil.valueToTree(config)); |
|||
when(ctxMock.getTimeseriesService()).thenReturn(timeseriesServiceMock); |
|||
|
|||
node.init(ctxMock, nodeConfiguration); |
|||
} |
|||
|
|||
@Test |
|||
public void givenDefaultConfig_whenDefaultConfiguration_thenVerify() { |
|||
assertEquals(config.getInputValueKey(), "pulseCounter"); |
|||
assertEquals(config.getOutputValueKey(), "delta"); |
|||
assertTrue(config.isUseCache()); |
|||
assertFalse(config.isAddPeriodBetweenMsgs()); |
|||
assertEquals(config.getPeriodValueKey(), "periodInMs"); |
|||
assertTrue(config.isTellFailureIfDeltaIsNegative()); |
|||
} |
|||
|
|||
@Test |
|||
public void givenInvalidMsgType_whenOnMsg_thenShouldTellNextOther() { |
|||
// GIVEN
|
|||
var msgData = "{\"pulseCounter\": 42}"; |
|||
var msg = TbMsg.newMsg("POST_ATTRIBUTES_REQUEST", DUMMY_DEVICE_ORIGINATOR, new TbMsgMetaData(), msgData); |
|||
|
|||
// WHEN
|
|||
node.onMsg(ctxMock, msg); |
|||
|
|||
// THEN
|
|||
verify(ctxMock, times(1)).tellNext(eq(msg), eq("Other")); |
|||
verify(ctxMock, never()).tellSuccess(any()); |
|||
verify(ctxMock, never()).tellFailure(any(), any()); |
|||
} |
|||
|
|||
@Test |
|||
public void givenInputKeyIsNotPresent_whenOnMsg_thenShouldTellNextOther() { |
|||
// GIVEN
|
|||
var msg = TbMsg.newMsg("POST_TELEMETRY_REQUEST", DUMMY_DEVICE_ORIGINATOR, new TbMsgMetaData(), "{}"); |
|||
|
|||
// WHEN
|
|||
node.onMsg(ctxMock, msg); |
|||
|
|||
// THEN
|
|||
verify(ctxMock, times(1)).tellNext(eq(msg), eq("Other")); |
|||
verify(ctxMock, never()).tellSuccess(any()); |
|||
verify(ctxMock, never()).tellFailure(any(), any()); |
|||
} |
|||
|
|||
@Test |
|||
public void givenDoubleValue_whenOnMsg_thenShouldTellSuccess() throws TbNodeException { |
|||
// GIVEN
|
|||
config.setRound(1); |
|||
config.setInputValueKey("temperature"); |
|||
config.setOutputValueKey("temp_delta"); |
|||
nodeConfiguration = new TbNodeConfiguration(JacksonUtil.valueToTree(config)); |
|||
node.init(ctxMock, nodeConfiguration); |
|||
|
|||
mockFindLatest(new BasicTsKvEntry(System.currentTimeMillis(), new DoubleDataEntry("temperature", 40.5))); |
|||
|
|||
var msgData = "{\"temperature\": 42,\"airPressure\":123}"; |
|||
var msg = TbMsg.newMsg("POST_TELEMETRY_REQUEST", DUMMY_DEVICE_ORIGINATOR, new TbMsgMetaData(), msgData); |
|||
|
|||
// WHEN
|
|||
node.onMsg(ctxMock, msg); |
|||
|
|||
// THEN
|
|||
var actualMsgCaptor = ArgumentCaptor.forClass(TbMsg.class); |
|||
|
|||
verify(ctxMock, times(1)).tellSuccess(actualMsgCaptor.capture()); |
|||
verify(ctxMock, never()).tellNext(any(), anyString()); |
|||
verify(ctxMock, never()).tellNext(any(), anySet()); |
|||
verify(ctxMock, never()).tellFailure(any(), any()); |
|||
|
|||
var expectedMsgData = "{\"temperature\":42,\"airPressure\":123,\"temp_delta\":1.5}"; |
|||
|
|||
assertEquals(expectedMsgData, actualMsgCaptor.getValue().getData()); |
|||
} |
|||
|
|||
@Test |
|||
public void givenLongStringValue_whenOnMsg_thenShouldTellSuccess() throws TbNodeException { |
|||
// GIVEN
|
|||
config.setInputValueKey("temperature"); |
|||
config.setOutputValueKey("temp_delta"); |
|||
nodeConfiguration = new TbNodeConfiguration(JacksonUtil.valueToTree(config)); |
|||
node.init(ctxMock, nodeConfiguration); |
|||
|
|||
mockFindLatest(new BasicTsKvEntry(System.currentTimeMillis(), new LongDataEntry("temperature", 40L))); |
|||
|
|||
var msgData = "{\"temperature\": 42,\"airPressure\":123}"; |
|||
var msg = TbMsg.newMsg("POST_TELEMETRY_REQUEST", DUMMY_DEVICE_ORIGINATOR, new TbMsgMetaData(), msgData); |
|||
|
|||
// WHEN
|
|||
node.onMsg(ctxMock, msg); |
|||
|
|||
// THEN
|
|||
var actualMsgCaptor = ArgumentCaptor.forClass(TbMsg.class); |
|||
|
|||
verify(ctxMock, times(1)).tellSuccess(actualMsgCaptor.capture()); |
|||
verify(ctxMock, never()).tellNext(any(), anyString()); |
|||
verify(ctxMock, never()).tellNext(any(), anySet()); |
|||
verify(ctxMock, never()).tellFailure(any(), any()); |
|||
|
|||
var expectedMsgData = "{\"temperature\":42,\"airPressure\":123,\"temp_delta\":2}"; |
|||
|
|||
assertEquals(expectedMsgData, actualMsgCaptor.getValue().getData()); |
|||
} |
|||
|
|||
@Test |
|||
public void givenValidStringValue_whenOnMsg_thenShouldTellSuccess() throws TbNodeException { |
|||
// GIVEN
|
|||
config.setInputValueKey("temperature"); |
|||
config.setOutputValueKey("temp_delta"); |
|||
nodeConfiguration = new TbNodeConfiguration(JacksonUtil.valueToTree(config)); |
|||
node.init(ctxMock, nodeConfiguration); |
|||
|
|||
mockFindLatest(new BasicTsKvEntry(System.currentTimeMillis(), new StringDataEntry("temperature", "40.0"))); |
|||
|
|||
var msgData = "{\"temperature\": 42,\"airPressure\":123}"; |
|||
var msg = TbMsg.newMsg("POST_TELEMETRY_REQUEST", DUMMY_DEVICE_ORIGINATOR, new TbMsgMetaData(), msgData); |
|||
|
|||
// WHEN
|
|||
node.onMsg(ctxMock, msg); |
|||
|
|||
// THEN
|
|||
var actualMsgCaptor = ArgumentCaptor.forClass(TbMsg.class); |
|||
|
|||
verify(ctxMock, times(1)).tellSuccess(actualMsgCaptor.capture()); |
|||
verify(ctxMock, never()).tellNext(any(), anyString()); |
|||
verify(ctxMock, never()).tellNext(any(), anySet()); |
|||
verify(ctxMock, never()).tellFailure(any(), any()); |
|||
|
|||
var expectedMsgData = "{\"temperature\":42,\"airPressure\":123,\"temp_delta\":2}"; |
|||
|
|||
assertEquals(expectedMsgData, actualMsgCaptor.getValue().getData()); |
|||
} |
|||
|
|||
@Test |
|||
public void givenTwoMessagesAndPeriodOnAndCachingOn_whenOnMsg_thenVerify() throws TbNodeException { |
|||
// STAGE 1
|
|||
// GIVEN
|
|||
config.setInputValueKey("temperature"); |
|||
config.setOutputValueKey("temp_delta"); |
|||
config.setPeriodValueKey("ts_delta"); |
|||
config.setAddPeriodBetweenMsgs(true); |
|||
nodeConfiguration = new TbNodeConfiguration(JacksonUtil.valueToTree(config)); |
|||
node.init(ctxMock, nodeConfiguration); |
|||
|
|||
mockFindLatest(new BasicTsKvEntry(1L, new DoubleDataEntry("temperature", 40.0))); |
|||
|
|||
var msgData = "{\"temperature\": 42,\"airPressure\":123}"; |
|||
var firstMsgMetaData = new TbMsgMetaData(); |
|||
firstMsgMetaData.putValue("ts", String.valueOf(3L)); |
|||
var firstMsg = TbMsg.newMsg("POST_TELEMETRY_REQUEST", DUMMY_DEVICE_ORIGINATOR, firstMsgMetaData, msgData); |
|||
|
|||
// WHEN
|
|||
node.onMsg(ctxMock, firstMsg); |
|||
|
|||
// THEN
|
|||
var actualMsgCaptor = ArgumentCaptor.forClass(TbMsg.class); |
|||
|
|||
verify(ctxMock, times(1)).tellSuccess(actualMsgCaptor.capture()); |
|||
verify(ctxMock, never()).tellNext(any(), anyString()); |
|||
verify(ctxMock, never()).tellNext(any(), anySet()); |
|||
verify(ctxMock, never()).tellFailure(any(), any()); |
|||
|
|||
var expectedMsgData = "{\"temperature\":42,\"airPressure\":123,\"temp_delta\":2,\"ts_delta\":2}"; |
|||
|
|||
assertEquals(expectedMsgData, actualMsgCaptor.getValue().getData()); |
|||
|
|||
// STAGE 2
|
|||
// GIVEN
|
|||
reset(ctxMock); |
|||
reset(timeseriesServiceMock); |
|||
|
|||
var secondMsgMetaData = new TbMsgMetaData(); |
|||
secondMsgMetaData.putValue("ts", String.valueOf(6L)); |
|||
var secondMsg = TbMsg.newMsg("POST_TELEMETRY_REQUEST", DUMMY_DEVICE_ORIGINATOR, secondMsgMetaData, msgData); |
|||
|
|||
// WHEN
|
|||
node.onMsg(ctxMock, secondMsg); |
|||
|
|||
// THEN
|
|||
actualMsgCaptor = ArgumentCaptor.forClass(TbMsg.class); |
|||
|
|||
verify(timeseriesServiceMock, never()).findLatest(any(), any(), anyList()); |
|||
verify(ctxMock, times(1)).tellSuccess(actualMsgCaptor.capture()); |
|||
verify(ctxMock, never()).tellNext(any(), anyString()); |
|||
verify(ctxMock, never()).tellNext(any(), anySet()); |
|||
verify(ctxMock, never()).tellFailure(any(), any()); |
|||
|
|||
expectedMsgData = "{\"temperature\":42,\"airPressure\":123,\"temp_delta\":0,\"ts_delta\":3}"; |
|||
|
|||
assertEquals(expectedMsgData, actualMsgCaptor.getValue().getData()); |
|||
} |
|||
|
|||
@Test |
|||
public void givenLastValueIsNull_whenOnMsh_thenDeltaShouldBeZero() throws TbNodeException { |
|||
// GIVEN
|
|||
config.setInputValueKey("temperature"); |
|||
config.setOutputValueKey("temp_delta"); |
|||
nodeConfiguration = new TbNodeConfiguration(JacksonUtil.valueToTree(config)); |
|||
node.init(ctxMock, nodeConfiguration); |
|||
|
|||
mockFindLatest(new BasicTsKvEntry(System.currentTimeMillis(), new DoubleDataEntry("temperature", null))); |
|||
|
|||
var msgData = "{\"temperature\": 42,\"airPressure\":123}"; |
|||
var msg = TbMsg.newMsg("POST_TELEMETRY_REQUEST", DUMMY_DEVICE_ORIGINATOR, new TbMsgMetaData(), msgData); |
|||
|
|||
// WHEN
|
|||
node.onMsg(ctxMock, msg); |
|||
|
|||
// THEN
|
|||
var actualMsgCaptor = ArgumentCaptor.forClass(TbMsg.class); |
|||
|
|||
verify(ctxMock, times(1)).tellSuccess(actualMsgCaptor.capture()); |
|||
verify(ctxMock, never()).tellNext(any(), anyString()); |
|||
verify(ctxMock, never()).tellNext(any(), anySet()); |
|||
verify(ctxMock, never()).tellFailure(any(), any()); |
|||
|
|||
var expectedMsgData = "{\"temperature\":42,\"airPressure\":123,\"temp_delta\":0}"; |
|||
|
|||
assertEquals(expectedMsgData, actualMsgCaptor.getValue().getData()); |
|||
} |
|||
|
|||
@Test |
|||
public void givenNegativeDeltaAndTellFailureIfNegativeDeltaTrue_whenOnMsg_thenShouldTellFailure() throws TbNodeException { |
|||
// GIVEN
|
|||
config.setTellFailureIfDeltaIsNegative(true); |
|||
nodeConfiguration = new TbNodeConfiguration(JacksonUtil.valueToTree(config)); |
|||
node.init(ctxMock, nodeConfiguration); |
|||
|
|||
mockFindLatest(new BasicTsKvEntry(System.currentTimeMillis(), new LongDataEntry("pulseCounter", 200L))); |
|||
|
|||
var msgData = "{\"pulseCounter\":\"123\"}"; |
|||
var msg = TbMsg.newMsg("POST_TELEMETRY_REQUEST", DUMMY_DEVICE_ORIGINATOR, new TbMsgMetaData(), msgData); |
|||
|
|||
// WHEN
|
|||
node.onMsg(ctxMock, msg); |
|||
|
|||
// THEN
|
|||
verify(ctxMock, times(1)).tellNext(msg, "Failure"); |
|||
verify(ctxMock, never()).tellSuccess(any()); |
|||
verify(ctxMock, never()).tellFailure(any(), any()); |
|||
verify(ctxMock, never()).tellNext(any(), anySet()); |
|||
} |
|||
|
|||
@Test |
|||
public void givenInvalidStringValue_whenOnMsg_thenException() { |
|||
// GIVEN
|
|||
mockFindLatest(new BasicTsKvEntry(System.currentTimeMillis(), new StringDataEntry("pulseCounter", "high"))); |
|||
|
|||
var msgData = "{\"pulseCounter\":\"123\"}"; |
|||
var msg = TbMsg.newMsg("POST_TELEMETRY_REQUEST", DUMMY_DEVICE_ORIGINATOR, new TbMsgMetaData(), msgData); |
|||
|
|||
// WHEN
|
|||
node.onMsg(ctxMock, msg); |
|||
|
|||
// THEN
|
|||
var actualMsgCaptor = ArgumentCaptor.forClass(TbMsg.class); |
|||
var actualExceptionCaptor = ArgumentCaptor.forClass(Exception.class); |
|||
|
|||
verify(ctxMock, times(1)).tellFailure(actualMsgCaptor.capture(), actualExceptionCaptor.capture()); |
|||
verify(ctxMock, never()).tellSuccess(any()); |
|||
verify(ctxMock, never()).tellNext(any(), anyString()); |
|||
verify(ctxMock, never()).tellNext(any(), anySet()); |
|||
|
|||
var expectedExceptionMsg = "Calculation failed. Unable to parse value [high] of telemetry [pulseCounter] to Double"; |
|||
var actualException = actualExceptionCaptor.getValue(); |
|||
|
|||
assertEquals(msg, actualMsgCaptor.getValue()); |
|||
assertInstanceOf(IllegalArgumentException.class, actualException); |
|||
assertEquals(expectedExceptionMsg, actualException.getMessage()); |
|||
} |
|||
|
|||
@Test |
|||
public void givenBooleanValue_whenOnMsg_thenException() { |
|||
// GIVEN
|
|||
mockFindLatest(new BasicTsKvEntry(System.currentTimeMillis(), new BooleanDataEntry("pulseCounter", false))); |
|||
|
|||
var msgData = "{\"pulseCounter\":true}"; |
|||
var msg = TbMsg.newMsg("POST_TELEMETRY_REQUEST", DUMMY_DEVICE_ORIGINATOR, new TbMsgMetaData(), msgData); |
|||
|
|||
// WHEN
|
|||
node.onMsg(ctxMock, msg); |
|||
|
|||
// THEN
|
|||
var actualMsgCaptor = ArgumentCaptor.forClass(TbMsg.class); |
|||
var actualExceptionCaptor = ArgumentCaptor.forClass(Exception.class); |
|||
|
|||
verify(ctxMock, times(1)).tellFailure(actualMsgCaptor.capture(), actualExceptionCaptor.capture()); |
|||
verify(ctxMock, never()).tellSuccess(any()); |
|||
verify(ctxMock, never()).tellNext(any(), anyString()); |
|||
verify(ctxMock, never()).tellNext(any(), anySet()); |
|||
|
|||
var expectedExceptionMsg = "Calculation failed. Boolean values are not supported!"; |
|||
var actualException = actualExceptionCaptor.getValue(); |
|||
|
|||
assertEquals(msg, actualMsgCaptor.getValue()); |
|||
assertInstanceOf(IllegalArgumentException.class, actualException); |
|||
assertEquals(expectedExceptionMsg, actualException.getMessage()); |
|||
} |
|||
|
|||
@Test |
|||
public void givenJsonValue_whenOnMsg_thenException() { |
|||
// GIVEN
|
|||
mockFindLatest(new BasicTsKvEntry(System.currentTimeMillis(), new JsonDataEntry("pulseCounter", "{\"isActive\":false}"))); |
|||
|
|||
var msgData = "{\"pulseCounter\":{\"isActive\":true}}"; |
|||
var msg = TbMsg.newMsg("POST_TELEMETRY_REQUEST", DUMMY_DEVICE_ORIGINATOR, new TbMsgMetaData(), msgData); |
|||
|
|||
// WHEN
|
|||
node.onMsg(ctxMock, msg); |
|||
|
|||
// THEN
|
|||
var actualMsgCaptor = ArgumentCaptor.forClass(TbMsg.class); |
|||
var actualExceptionCaptor = ArgumentCaptor.forClass(Exception.class); |
|||
|
|||
verify(ctxMock, times(1)).tellFailure(actualMsgCaptor.capture(), actualExceptionCaptor.capture()); |
|||
verify(ctxMock, never()).tellSuccess(any()); |
|||
verify(ctxMock, never()).tellNext(any(), anyString()); |
|||
verify(ctxMock, never()).tellNext(any(), anySet()); |
|||
|
|||
var expectedExceptionMsg = "Calculation failed. JSON values are not supported!"; |
|||
var actualException = actualExceptionCaptor.getValue(); |
|||
|
|||
assertEquals(msg, actualMsgCaptor.getValue()); |
|||
assertInstanceOf(IllegalArgumentException.class, actualException); |
|||
assertEquals(expectedExceptionMsg, actualException.getMessage()); |
|||
} |
|||
|
|||
private void mockFindLatest(TsKvEntry tsKvEntry) { |
|||
when(ctxMock.getDbCallbackExecutor()).thenReturn(DB_EXECUTOR); |
|||
when(ctxMock.getTenantId()).thenReturn(TENANT_ID); |
|||
when(timeseriesServiceMock.findLatest( |
|||
eq(TENANT_ID), eq(DUMMY_DEVICE_ORIGINATOR), argThat(new ListMatcher<>(List.of(tsKvEntry.getKey()))) |
|||
)).thenReturn(Futures.immediateFuture(List.of(tsKvEntry))); |
|||
} |
|||
|
|||
@RequiredArgsConstructor |
|||
private static class ListMatcher<T> implements ArgumentMatcher<List<T>> { |
|||
|
|||
private final List<T> expectedList; |
|||
|
|||
@Override |
|||
public boolean matches(List<T> actualList) { |
|||
if (actualList == expectedList) { |
|||
return true; |
|||
} |
|||
if (actualList.size() != expectedList.size()) { |
|||
return false; |
|||
} |
|||
return actualList.containsAll(expectedList); |
|||
} |
|||
|
|||
} |
|||
|
|||
} |
|||
@ -1,238 +0,0 @@ |
|||
/** |
|||
* Copyright © 2016-2023 The Thingsboard Authors |
|||
* |
|||
* Licensed under the Apache License, Version 2.0 (the "License"); |
|||
* you may not use this file except in compliance with the License. |
|||
* You may obtain a copy of the License at |
|||
* |
|||
* http://www.apache.org/licenses/LICENSE-2.0
|
|||
* |
|||
* Unless required by applicable law or agreed to in writing, software |
|||
* distributed under the License is distributed on an "AS IS" BASIS, |
|||
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. |
|||
* See the License for the specific language governing permissions and |
|||
* limitations under the License. |
|||
*/ |
|||
package org.thingsboard.rule.engine.metadata; |
|||
|
|||
import com.datastax.oss.driver.api.core.uuid.Uuids; |
|||
import com.fasterxml.jackson.databind.ObjectMapper; |
|||
import com.google.common.collect.Lists; |
|||
import com.google.common.util.concurrent.Futures; |
|||
import org.junit.runner.RunWith; |
|||
import org.mockito.ArgumentCaptor; |
|||
import org.mockito.Mock; |
|||
import org.mockito.junit.MockitoJUnitRunner; |
|||
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.server.common.data.Device; |
|||
import org.thingsboard.server.common.data.User; |
|||
import org.thingsboard.server.common.data.asset.Asset; |
|||
import org.thingsboard.server.common.data.id.CustomerId; |
|||
import org.thingsboard.server.common.data.id.EntityId; |
|||
import org.thingsboard.server.common.data.id.RuleChainId; |
|||
import org.thingsboard.server.common.data.id.RuleNodeId; |
|||
import org.thingsboard.server.common.data.id.TenantId; |
|||
import org.thingsboard.server.common.data.kv.AttributeKvEntry; |
|||
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.TsKvEntry; |
|||
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.dao.asset.AssetService; |
|||
import org.thingsboard.server.dao.attributes.AttributesService; |
|||
import org.thingsboard.server.dao.device.DeviceService; |
|||
import org.thingsboard.server.dao.timeseries.TimeseriesService; |
|||
import org.thingsboard.server.dao.user.UserService; |
|||
|
|||
import java.util.HashMap; |
|||
import java.util.List; |
|||
import java.util.Map; |
|||
import java.util.NoSuchElementException; |
|||
|
|||
import static org.assertj.core.api.Assertions.assertThat; |
|||
import static org.junit.Assert.assertEquals; |
|||
import static org.junit.Assert.assertTrue; |
|||
import static org.mockito.ArgumentMatchers.any; |
|||
import static org.mockito.ArgumentMatchers.anyCollection; |
|||
import static org.mockito.ArgumentMatchers.eq; |
|||
import static org.mockito.ArgumentMatchers.same; |
|||
import static org.mockito.Mockito.verify; |
|||
import static org.mockito.Mockito.when; |
|||
import static org.thingsboard.server.common.data.DataConstants.SERVER_SCOPE; |
|||
|
|||
@RunWith(MockitoJUnitRunner.class) |
|||
public abstract class TbAbstractAttributeNodeTest { |
|||
final CustomerId customerId = new CustomerId(Uuids.timeBased()); |
|||
final TenantId tenantId = TenantId.fromUUID(Uuids.timeBased()); |
|||
final RuleChainId ruleChainId = new RuleChainId(Uuids.timeBased()); |
|||
final RuleNodeId ruleNodeId = new RuleNodeId(Uuids.timeBased()); |
|||
final String keyAttrConf = "${word}"; |
|||
final String valueAttrConf = "${result}"; |
|||
@Mock |
|||
TbContext ctx; |
|||
@Mock |
|||
AttributesService attributesService; |
|||
@Mock |
|||
TimeseriesService timeseriesService; |
|||
@Mock |
|||
UserService userService; |
|||
@Mock |
|||
AssetService assetService; |
|||
@Mock |
|||
DeviceService deviceService; |
|||
TbMsg msg; |
|||
Map<String, String> metaData; |
|||
TbAbstractGetEntityAttrNode node; |
|||
|
|||
void init(TbAbstractGetEntityAttrNode node) throws TbNodeException { |
|||
ObjectMapper mapper = JacksonUtil.OBJECT_MAPPER; |
|||
TbNodeConfiguration nodeConfiguration = new TbNodeConfiguration(mapper.valueToTree(getTbNodeConfig())); |
|||
|
|||
metaData = new HashMap<>(); |
|||
metaData.putIfAbsent("word", "temperature"); |
|||
metaData.putIfAbsent("result", "answer"); |
|||
|
|||
this.node = node; |
|||
this.node.init(null, nodeConfiguration); |
|||
} |
|||
|
|||
void errorThrownIfCannotLoadAttributes(User user) { |
|||
msg = TbMsg.newMsg("USER", user.getId(), new TbMsgMetaData(), TbMsgDataType.JSON, "{}", ruleChainId, ruleNodeId); |
|||
|
|||
when(ctx.getAttributesService()).thenReturn(attributesService); |
|||
when(attributesService.find(any(), eq(getEntityId()), eq(SERVER_SCOPE), anyCollection())) |
|||
.thenThrow(new IllegalStateException("something wrong")); |
|||
|
|||
node.onMsg(ctx, msg); |
|||
final ArgumentCaptor<Throwable> captor = ArgumentCaptor.forClass(Throwable.class); |
|||
verify(ctx).tellFailure(same(msg), captor.capture()); |
|||
|
|||
Throwable value = captor.getValue(); |
|||
assertEquals("something wrong", value.getMessage()); |
|||
assertTrue(msg.getMetaData().getData().isEmpty()); |
|||
} |
|||
|
|||
void errorThrownIfCannotLoadAttributesAsync(User user) { |
|||
msg = TbMsg.newMsg("USER", user.getId(), new TbMsgMetaData(), TbMsgDataType.JSON, "{}", ruleChainId, ruleNodeId); |
|||
|
|||
when(ctx.getAttributesService()).thenReturn(attributesService); |
|||
when(attributesService.find(any(), eq(getEntityId()), eq(SERVER_SCOPE), anyCollection())) |
|||
.thenReturn(Futures.immediateFailedFuture(new IllegalStateException("something wrong"))); |
|||
|
|||
node.onMsg(ctx, msg); |
|||
final ArgumentCaptor<Throwable> captor = ArgumentCaptor.forClass(Throwable.class); |
|||
verify(ctx).tellFailure(same(msg), captor.capture()); |
|||
|
|||
Throwable value = captor.getValue(); |
|||
assertEquals("something wrong", value.getMessage()); |
|||
assertTrue(msg.getMetaData().getData().isEmpty()); |
|||
} |
|||
|
|||
void failedChainUsedIfCustomerCannotBeFound(User user) { |
|||
msg = TbMsg.newMsg("USER", user.getId(), new TbMsgMetaData(), TbMsgDataType.JSON, "{}", ruleChainId, ruleNodeId); |
|||
|
|||
node.onMsg(ctx, msg); |
|||
var exceptionCaptor = ArgumentCaptor.forClass(NoSuchElementException.class); |
|||
verify(ctx).tellFailure(eq(msg), exceptionCaptor.capture()); |
|||
|
|||
assertThat(exceptionCaptor.getValue().getMessage()).contains("Did not find entity! Msg ID: "); |
|||
assertTrue(msg.getMetaData().getData().isEmpty()); |
|||
} |
|||
|
|||
void entityAttributeAddedInMetadata(EntityId entityId, String type) { |
|||
msg = TbMsg.newMsg(type, entityId, new TbMsgMetaData(metaData), TbMsgDataType.JSON, "{}", ruleChainId, ruleNodeId); |
|||
entityAttributeFetched(getEntityId()); |
|||
} |
|||
|
|||
void usersCustomerAttributesFetched(User user) { |
|||
msg = TbMsg.newMsg("USER", user.getId(), new TbMsgMetaData(metaData), TbMsgDataType.JSON, "{}", ruleChainId, ruleNodeId); |
|||
|
|||
entityAttributeFetched(getEntityId()); |
|||
} |
|||
|
|||
void assetsCustomerAttributesFetched(Asset asset) { |
|||
msg = TbMsg.newMsg("ASSET", asset.getId(), new TbMsgMetaData(metaData), TbMsgDataType.JSON, "{}", ruleChainId, ruleNodeId); |
|||
|
|||
entityAttributeFetched(getEntityId()); |
|||
} |
|||
|
|||
void deviceCustomerAttributesFetched(Device device) { |
|||
msg = TbMsg.newMsg("DEVICE", device.getId(), new TbMsgMetaData(metaData), TbMsgDataType.JSON, "{}", ruleChainId, ruleNodeId); |
|||
|
|||
entityAttributeFetched(getEntityId()); |
|||
} |
|||
|
|||
void deviceCustomerTelemetryFetched(Device device) throws TbNodeException { |
|||
ObjectMapper mapper = JacksonUtil.OBJECT_MAPPER; |
|||
TbNodeConfiguration nodeConfiguration = new TbNodeConfiguration(mapper.valueToTree(getTbNodeConfigForTelemetry())); |
|||
|
|||
TbAbstractGetEntityAttrNode node = getEmptyNode(); |
|||
node.init(null, nodeConfiguration); |
|||
|
|||
msg = TbMsg.newMsg("DEVICE", device.getId(), new TbMsgMetaData(metaData), TbMsgDataType.JSON, "{}", ruleChainId, ruleNodeId); |
|||
|
|||
List<TsKvEntry> timeseries = Lists.newArrayList(new BasicTsKvEntry(1L, new StringDataEntry("temperature", "highest"))); |
|||
|
|||
when(ctx.getTimeseriesService()).thenReturn(timeseriesService); |
|||
when(timeseriesService.findLatest(any(), eq(getEntityId()), anyCollection())) |
|||
.thenReturn(Futures.immediateFuture(timeseries)); |
|||
|
|||
node.onMsg(ctx, msg); |
|||
verify(ctx).tellSuccess(msg); |
|||
assertEquals(msg.getMetaData().getValue("answer"), "highest"); |
|||
} |
|||
|
|||
void entityAttributeFetched(EntityId entityId) { |
|||
List<AttributeKvEntry> attributes = Lists.newArrayList(new BaseAttributeKvEntry(new StringDataEntry("temperature", "high"), 1L)); |
|||
|
|||
when(ctx.getAttributesService()).thenReturn(attributesService); |
|||
when(attributesService.find(any(), eq(entityId), eq(SERVER_SCOPE), anyCollection())) |
|||
.thenReturn(Futures.immediateFuture(attributes)); |
|||
|
|||
node.onMsg(ctx, msg); |
|||
verify(ctx).tellSuccess(msg); |
|||
assertEquals(msg.getMetaData().getValue("answer"), "high"); |
|||
} |
|||
|
|||
TbGetEntityAttrNodeConfiguration getTbNodeConfig() { |
|||
return getConfig(false); |
|||
} |
|||
|
|||
TbGetEntityAttrNodeConfiguration getTbNodeConfigForTelemetry() { |
|||
return getConfig(true); |
|||
} |
|||
|
|||
private TbGetEntityAttrNodeConfiguration getConfig(boolean isTelemetry) { |
|||
TbGetEntityAttrNodeConfiguration config = new TbGetEntityAttrNodeConfiguration(); |
|||
Map<String, String> conf = new HashMap<>(); |
|||
conf.put(keyAttrConf, valueAttrConf); |
|||
config.setAttrMapping(conf); |
|||
config.setTelemetry(isTelemetry); |
|||
config.setFetchTo(FetchTo.METADATA); |
|||
return config; |
|||
} |
|||
|
|||
protected abstract TbAbstractGetEntityAttrNode getEmptyNode(); |
|||
|
|||
abstract EntityId getEntityId(); |
|||
|
|||
void mockFindDevice(Device device) { |
|||
when(ctx.getDeviceService()).thenReturn(deviceService); |
|||
when(deviceService.findDeviceByIdAsync(any(), eq(device.getId()))).thenReturn(Futures.immediateFuture(device)); |
|||
} |
|||
|
|||
void mockFindAsset(Asset asset) { |
|||
when(ctx.getAssetService()).thenReturn(assetService); |
|||
when(assetService.findAssetByIdAsync(any(), eq(asset.getId()))).thenReturn(Futures.immediateFuture(asset)); |
|||
} |
|||
|
|||
void mockFindUser(User user) { |
|||
when(ctx.getUserService()).thenReturn(userService); |
|||
when(userService.findUserByIdAsync(any(), eq(user.getId()))).thenReturn(Futures.immediateFuture(user)); |
|||
} |
|||
} |
|||
@ -0,0 +1,469 @@ |
|||
/** |
|||
* Copyright © 2016-2023 The Thingsboard Authors |
|||
* |
|||
* Licensed under the Apache License, Version 2.0 (the "License"); |
|||
* you may not use this file except in compliance with the License. |
|||
* You may obtain a copy of the License at |
|||
* |
|||
* http://www.apache.org/licenses/LICENSE-2.0
|
|||
* |
|||
* Unless required by applicable law or agreed to in writing, software |
|||
* distributed under the License is distributed on an "AS IS" BASIS, |
|||
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. |
|||
* See the License for the specific language governing permissions and |
|||
* limitations under the License. |
|||
*/ |
|||
package org.thingsboard.rule.engine.metadata; |
|||
|
|||
import com.google.common.util.concurrent.Futures; |
|||
import com.google.common.util.concurrent.ListenableFuture; |
|||
import org.jetbrains.annotations.NotNull; |
|||
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.Mock; |
|||
import org.mockito.junit.jupiter.MockitoExtension; |
|||
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.TbNodeException; |
|||
import org.thingsboard.rule.engine.util.EntityDetails; |
|||
import org.thingsboard.server.common.data.Customer; |
|||
import org.thingsboard.server.common.data.Dashboard; |
|||
import org.thingsboard.server.common.data.Device; |
|||
import org.thingsboard.server.common.data.User; |
|||
import org.thingsboard.server.common.data.asset.Asset; |
|||
import org.thingsboard.server.common.data.edge.Edge; |
|||
import org.thingsboard.server.common.data.id.AssetId; |
|||
import org.thingsboard.server.common.data.id.CustomerId; |
|||
import org.thingsboard.server.common.data.id.DashboardId; |
|||
import org.thingsboard.server.common.data.id.DeviceId; |
|||
import org.thingsboard.server.common.data.id.EdgeId; |
|||
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.msg.TbMsg; |
|||
import org.thingsboard.server.common.msg.TbMsgMetaData; |
|||
import org.thingsboard.server.dao.asset.AssetService; |
|||
import org.thingsboard.server.dao.customer.CustomerService; |
|||
import org.thingsboard.server.dao.device.DeviceService; |
|||
import org.thingsboard.server.dao.edge.EdgeService; |
|||
import org.thingsboard.server.dao.entityview.EntityViewService; |
|||
import org.thingsboard.server.dao.user.UserService; |
|||
|
|||
import java.util.Collections; |
|||
import java.util.List; |
|||
import java.util.NoSuchElementException; |
|||
import java.util.UUID; |
|||
import java.util.concurrent.Callable; |
|||
|
|||
import static org.assertj.core.api.Assertions.assertThat; |
|||
import static org.junit.jupiter.api.Assertions.assertThrows; |
|||
import static org.mockito.ArgumentMatchers.any; |
|||
import static org.mockito.ArgumentMatchers.eq; |
|||
import static org.mockito.Mockito.never; |
|||
import static org.mockito.Mockito.times; |
|||
import static org.mockito.Mockito.verify; |
|||
import static org.mockito.Mockito.when; |
|||
|
|||
@ExtendWith(MockitoExtension.class) |
|||
public class TbGetCustomerDetailsNodeTest { |
|||
|
|||
private static final DeviceId DUMMY_DEVICE_ORIGINATOR = new DeviceId(UUID.randomUUID()); |
|||
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 |
|||
private TbContext ctxMock; |
|||
@Mock |
|||
private CustomerService customerServiceMock; |
|||
@Mock |
|||
private DeviceService deviceServiceMock; |
|||
@Mock |
|||
private AssetService assetServiceMock; |
|||
@Mock |
|||
private EntityViewService entityViewServiceMock; |
|||
@Mock |
|||
private UserService userServiceMock; |
|||
@Mock |
|||
private EdgeService edgeServiceMock; |
|||
private TbGetCustomerDetailsNode node; |
|||
private TbGetCustomerDetailsNodeConfiguration config; |
|||
private TbNodeConfiguration nodeConfiguration; |
|||
private TbMsg msg; |
|||
private Customer customer; |
|||
|
|||
@BeforeEach |
|||
public void setUp() { |
|||
node = new TbGetCustomerDetailsNode(); |
|||
config = new TbGetCustomerDetailsNodeConfiguration().defaultConfiguration(); |
|||
nodeConfiguration = new TbNodeConfiguration(JacksonUtil.valueToTree(config)); |
|||
customer = new Customer(); |
|||
customer.setId(new CustomerId(UUID.randomUUID())); |
|||
customer.setTitle("Customer title"); |
|||
customer.setCountry("Customer country"); |
|||
customer.setCity("Customer city"); |
|||
customer.setState("Customer state"); |
|||
customer.setZip("123456"); |
|||
customer.setAddress("Customer address 1"); |
|||
customer.setAddress2("Customer address 2"); |
|||
customer.setPhone("+123456789"); |
|||
customer.setEmail("email@tenant.com"); |
|||
customer.setAdditionalInfo(JacksonUtil.toJsonNode("{\"someProperty\":\"someValue\",\"description\":\"Customer description\"}")); |
|||
} |
|||
|
|||
@Test |
|||
public void givenConfigWithNullFetchTo_whenInit_thenException() { |
|||
// GIVEN
|
|||
config.setDetailsList(List.of(EntityDetails.ID)); |
|||
config.setFetchTo(null); |
|||
nodeConfiguration = new TbNodeConfiguration(JacksonUtil.valueToTree(config)); |
|||
|
|||
// WHEN
|
|||
var exception = assertThrows(TbNodeException.class, () -> node.init(ctxMock, nodeConfiguration)); |
|||
|
|||
// THEN
|
|||
assertThat(exception.getMessage()).isEqualTo("FetchTo cannot be null!"); |
|||
verify(ctxMock, never()).tellSuccess(any()); |
|||
} |
|||
|
|||
@Test |
|||
public void givenDefaultConfig_whenInit_thenOK() { |
|||
assertThat(config.getDetailsList()).isEqualTo(Collections.emptyList()); |
|||
assertThat(config.getFetchTo()).isEqualTo(FetchTo.DATA); |
|||
} |
|||
|
|||
@Test |
|||
public void givenCustomConfig_whenInit_thenOK() throws TbNodeException { |
|||
// GIVEN
|
|||
config.setDetailsList(List.of(EntityDetails.ID, EntityDetails.PHONE)); |
|||
config.setFetchTo(FetchTo.METADATA); |
|||
nodeConfiguration = new TbNodeConfiguration(JacksonUtil.valueToTree(config)); |
|||
|
|||
// WHEN
|
|||
node.init(ctxMock, nodeConfiguration); |
|||
|
|||
// THEN
|
|||
assertThat(node.config).isEqualTo(config); |
|||
assertThat(config.getDetailsList()).isEqualTo(List.of(EntityDetails.ID, EntityDetails.PHONE)); |
|||
assertThat(config.getFetchTo()).isEqualTo(FetchTo.METADATA); |
|||
assertThat(node.fetchTo).isEqualTo(FetchTo.METADATA); |
|||
} |
|||
|
|||
@Test |
|||
public void givenMsgDataIsNotAnJsonObjectAndFetchToData_whenOnMsg_thenException() { |
|||
// GIVEN
|
|||
node.fetchTo = FetchTo.DATA; |
|||
msg = TbMsg.newMsg("SOME_MESSAGE_TYPE", DUMMY_DEVICE_ORIGINATOR, new TbMsgMetaData(), "[]"); |
|||
|
|||
// WHEN
|
|||
var exception = assertThrows(IllegalArgumentException.class, () -> node.onMsg(ctxMock, msg)); |
|||
|
|||
// THEN
|
|||
assertThat(exception.getMessage()).isEqualTo("Message body is not an object!"); |
|||
verify(ctxMock, never()).tellSuccess(any()); |
|||
} |
|||
|
|||
@Test |
|||
public void givenAllEntityDetailsAndFetchToData_whenOnMsg_thenShouldTellSuccessAndFetchAllToData() { |
|||
// GIVEN
|
|||
var device = new Device(); |
|||
device.setId(new DeviceId(UUID.randomUUID())); |
|||
device.setCustomerId(customer.getId()); |
|||
|
|||
prepareMsgAndConfig(FetchTo.DATA, List.of(EntityDetails.values()), device.getId()); |
|||
|
|||
when(ctxMock.getDeviceService()).thenReturn(deviceServiceMock); |
|||
when(deviceServiceMock.findDeviceByIdAsync(eq(TENANT_ID), eq(device.getId()))).thenReturn(Futures.immediateFuture(device)); |
|||
|
|||
mockFindCustomer(); |
|||
|
|||
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 = "{\"dataKey1\":123,\"dataKey2\":\"dataValue2\"," + |
|||
"\"customer_id\":\"" + customer.getId() + "\"," + |
|||
"\"customer_title\":\"" + customer.getTitle() + "\"," + |
|||
"\"customer_country\":\"" + customer.getCountry() + "\"," + |
|||
"\"customer_city\":\"" + customer.getCity() + "\"," + |
|||
"\"customer_state\":\"" + customer.getState() + "\"," + |
|||
"\"customer_zip\":\"" + customer.getZip() + "\"," + |
|||
"\"customer_address\":\"" + customer.getAddress() + "\"," + |
|||
"\"customer_address2\":\"" + customer.getAddress2() + "\"," + |
|||
"\"customer_phone\":\"" + customer.getPhone() + "\"," + |
|||
"\"customer_email\":\"" + customer.getEmail() + "\"," + |
|||
"\"customer_additionalInfo\":\"" + customer.getAdditionalInfo().get("description").asText() + "\"}"; |
|||
|
|||
assertThat(actualMessageCaptor.getValue().getData()).isEqualTo(expectedMsgData); |
|||
assertThat(actualMessageCaptor.getValue().getMetaData()).isEqualTo(msg.getMetaData()); |
|||
} |
|||
|
|||
@Test |
|||
public void givenSomeEntityDetailsAndFetchToMetadata_whenOnMsg_thenShouldTellSuccessAndFetchSomeToMetaData() { |
|||
// GIVEN
|
|||
var asset = new Asset(); |
|||
asset.setId(new AssetId(UUID.randomUUID())); |
|||
asset.setCustomerId(customer.getId()); |
|||
|
|||
prepareMsgAndConfig(FetchTo.METADATA, List.of(EntityDetails.ID, EntityDetails.TITLE, EntityDetails.PHONE), asset.getId()); |
|||
|
|||
when(ctxMock.getAssetService()).thenReturn(assetServiceMock); |
|||
when(assetServiceMock.findAssetByIdAsync(eq(TENANT_ID), eq(asset.getId()))).thenReturn(Futures.immediateFuture(asset)); |
|||
|
|||
mockFindCustomer(); |
|||
|
|||
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(msg.getMetaData().getData()); |
|||
expectedMsgMetaData.putValue("customer_id", customer.getId().getId().toString()); |
|||
expectedMsgMetaData.putValue("customer_title", customer.getTitle()); |
|||
expectedMsgMetaData.putValue("customer_phone", customer.getPhone()); |
|||
|
|||
assertThat(actualMessageCaptor.getValue().getData()).isEqualTo(msg.getData()); |
|||
assertThat(actualMessageCaptor.getValue().getMetaData()).isEqualTo(expectedMsgMetaData); |
|||
} |
|||
|
|||
@Test |
|||
public void givenNotPresentEntityDetailsAndFetchToData_whenOnMsg_thenShouldTellSuccessAndFetchNothingToData() { |
|||
// GIVEN
|
|||
customer.setZip(null); |
|||
customer.setAddress(null); |
|||
customer.setAddress2(null); |
|||
|
|||
var user = new User(); |
|||
user.setId(new UserId(UUID.randomUUID())); |
|||
user.setCustomerId(customer.getId()); |
|||
|
|||
prepareMsgAndConfig(FetchTo.DATA, List.of(EntityDetails.ZIP, EntityDetails.ADDRESS, EntityDetails.ADDRESS2), user.getId()); |
|||
|
|||
when(ctxMock.getUserService()).thenReturn(userServiceMock); |
|||
when(userServiceMock.findUserByIdAsync(eq(TENANT_ID), eq(user.getId()))).thenReturn(Futures.immediateFuture(user)); |
|||
|
|||
when(ctxMock.getDbCallbackExecutor()).thenReturn(DB_EXECUTOR); |
|||
|
|||
mockFindCustomer(); |
|||
|
|||
// 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()); |
|||
|
|||
assertThat(actualMessageCaptor.getValue().getData()).isEqualTo(msg.getData()); |
|||
assertThat(actualMessageCaptor.getValue().getMetaData()).isEqualTo(msg.getMetaData()); |
|||
} |
|||
|
|||
@Test |
|||
public void givenDidNotFindCustomer_whenOnMsg_thenShouldTellSuccessAndFetchNothingToData() { |
|||
// GIVEN
|
|||
var edge = new Edge(); |
|||
edge.setId(new EdgeId(UUID.randomUUID())); |
|||
edge.setCustomerId(customer.getId()); |
|||
|
|||
prepareMsgAndConfig(FetchTo.DATA, List.of(EntityDetails.ZIP, EntityDetails.ADDRESS, EntityDetails.ADDRESS2), edge.getId()); |
|||
|
|||
when(ctxMock.getTenantId()).thenReturn(TENANT_ID); |
|||
|
|||
when(ctxMock.getEdgeService()).thenReturn(edgeServiceMock); |
|||
when(edgeServiceMock.findEdgeByIdAsync(eq(TENANT_ID), eq(edge.getId()))).thenReturn(Futures.immediateFuture(edge)); |
|||
|
|||
when(ctxMock.getCustomerService()).thenReturn(customerServiceMock); |
|||
when(customerServiceMock.findCustomerByIdAsync(eq(TENANT_ID), eq(customer.getId()))).thenReturn(Futures.immediateFuture(null)); |
|||
|
|||
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()); |
|||
|
|||
assertThat(actualMessageCaptor.getValue().getData()).isEqualTo(msg.getData()); |
|||
assertThat(actualMessageCaptor.getValue().getMetaData()).isEqualTo(msg.getMetaData()); |
|||
} |
|||
|
|||
@Test |
|||
public void givenDidNotFindOriginator_whenOnMsg_thenShouldTellSuccessAndFetchNothingToData() { |
|||
// GIVEN
|
|||
var edge = new Edge(); |
|||
edge.setId(new EdgeId(UUID.randomUUID())); |
|||
edge.setCustomerId(customer.getId()); |
|||
|
|||
prepareMsgAndConfig(FetchTo.DATA, List.of(EntityDetails.ZIP, EntityDetails.ADDRESS, EntityDetails.ADDRESS2), edge.getId()); |
|||
|
|||
when(ctxMock.getTenantId()).thenReturn(TENANT_ID); |
|||
|
|||
when(ctxMock.getEdgeService()).thenReturn(edgeServiceMock); |
|||
when(edgeServiceMock.findEdgeByIdAsync(eq(TENANT_ID), eq(edge.getId()))).thenReturn(Futures.immediateFuture(null)); |
|||
|
|||
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()); |
|||
|
|||
assertThat(actualMessageCaptor.getValue().getData()).isEqualTo(msg.getData()); |
|||
assertThat(actualMessageCaptor.getValue().getMetaData()).isEqualTo(msg.getMetaData()); |
|||
} |
|||
|
|||
@Test |
|||
public void givenOriginatorNotAssignedToCustomer_whenOnMsg_thenShouldTellFailureAndFetchNothingToData() { |
|||
// GIVEN
|
|||
var device = new Device(); |
|||
device.setId(new DeviceId(UUID.randomUUID())); |
|||
device.setName("Thermostat"); |
|||
|
|||
prepareMsgAndConfig(FetchTo.DATA, List.of(EntityDetails.ZIP, EntityDetails.ADDRESS, EntityDetails.ADDRESS2), device.getId()); |
|||
|
|||
when(ctxMock.getTenantId()).thenReturn(TENANT_ID); |
|||
|
|||
when(ctxMock.getDeviceService()).thenReturn(deviceServiceMock); |
|||
when(deviceServiceMock.findDeviceByIdAsync(eq(TENANT_ID), eq(device.getId()))).thenReturn(Futures.immediateFuture(device)); |
|||
|
|||
when(ctxMock.getDbCallbackExecutor()).thenReturn(DB_EXECUTOR); |
|||
|
|||
// WHEN
|
|||
node.onMsg(ctxMock, msg); |
|||
|
|||
// THEN
|
|||
var actualMessageCaptor = ArgumentCaptor.forClass(TbMsg.class); |
|||
var actualExceptionCaptor = ArgumentCaptor.forClass(Exception.class); |
|||
|
|||
verify(ctxMock, times(1)).tellFailure(actualMessageCaptor.capture(), actualExceptionCaptor.capture()); |
|||
verify(ctxMock, never()).tellSuccess(any()); |
|||
|
|||
var actualMsg = actualMessageCaptor.getValue(); |
|||
var actualException = actualExceptionCaptor.getValue(); |
|||
|
|||
assertThat(actualMsg.getData()).isEqualTo(msg.getData()); |
|||
assertThat(actualMsg.getMetaData()).isEqualTo(msg.getMetaData()); |
|||
|
|||
assertThat(actualException).isInstanceOf(RuntimeException.class); |
|||
assertThat(actualException.getMessage()).isEqualTo("Device with name 'Thermostat' is not assigned to Customer."); |
|||
} |
|||
|
|||
@Test |
|||
public void givenNullDescriptionAndAddInfoEntityDetails_whenOnMsg_thenShouldTellSuccessAndFetchNothingToData() { |
|||
// GIVEN
|
|||
customer.setAdditionalInfo(JacksonUtil.toJsonNode("{\"someProperty\":\"someValue\",\"description\":null}")); |
|||
|
|||
var device = new Device(); |
|||
device.setId(new DeviceId(UUID.randomUUID())); |
|||
device.setCustomerId(customer.getId()); |
|||
|
|||
prepareMsgAndConfig(FetchTo.DATA, List.of(EntityDetails.ADDITIONAL_INFO), device.getId()); |
|||
|
|||
when(ctxMock.getDeviceService()).thenReturn(deviceServiceMock); |
|||
when(deviceServiceMock.findDeviceByIdAsync(eq(TENANT_ID), eq(device.getId()))).thenReturn(Futures.immediateFuture(device)); |
|||
|
|||
when(ctxMock.getDbCallbackExecutor()).thenReturn(DB_EXECUTOR); |
|||
|
|||
mockFindCustomer(); |
|||
|
|||
// 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()); |
|||
|
|||
assertThat(actualMessageCaptor.getValue().getData()).isEqualTo(msg.getData()); |
|||
assertThat(actualMessageCaptor.getValue().getMetaData()).isEqualTo(msg.getMetaData()); |
|||
} |
|||
|
|||
@Test |
|||
public void givenUnsupportedEntityType_whenOnMsg_thenShouldTellFailureAndFetchNothingToMetaData() { |
|||
// GIVEN
|
|||
var dashboard = new Dashboard(); |
|||
dashboard.setId(new DashboardId(UUID.randomUUID())); |
|||
|
|||
prepareMsgAndConfig(FetchTo.METADATA, List.of(EntityDetails.STATE), dashboard.getId()); |
|||
|
|||
// WHEN
|
|||
node.onMsg(ctxMock, msg); |
|||
|
|||
// THEN
|
|||
var actualMessageCaptor = ArgumentCaptor.forClass(TbMsg.class); |
|||
var actualExceptionCaptor = ArgumentCaptor.forClass(Exception.class); |
|||
|
|||
verify(ctxMock, times(1)).tellFailure(actualMessageCaptor.capture(), actualExceptionCaptor.capture()); |
|||
verify(ctxMock, never()).tellSuccess(any()); |
|||
|
|||
var actualMsg = actualMessageCaptor.getValue(); |
|||
var actualException = actualExceptionCaptor.getValue(); |
|||
|
|||
assertThat(actualMsg.getData()).isEqualTo(msg.getData()); |
|||
assertThat(actualMsg.getMetaData()).isEqualTo(msg.getMetaData()); |
|||
|
|||
assertThat(actualException).isInstanceOf(NoSuchElementException.class); |
|||
assertThat(actualException.getMessage()).isEqualTo("Entity with entityType 'DASHBOARD' is not supported."); |
|||
} |
|||
|
|||
private void prepareMsgAndConfig(FetchTo fetchTo, List<EntityDetails> detailsList, EntityId originator) { |
|||
config.setDetailsList(detailsList); |
|||
config.setFetchTo(fetchTo); |
|||
|
|||
node.config = config; |
|||
node.fetchTo = fetchTo; |
|||
|
|||
var msgMetaData = new TbMsgMetaData(); |
|||
msgMetaData.putValue("metaKey1", "metaValue1"); |
|||
msgMetaData.putValue("metaKey2", "metaValue2"); |
|||
|
|||
var msgData = "{\"dataKey1\":123,\"dataKey2\":\"dataValue2\"}"; |
|||
|
|||
msg = TbMsg.newMsg("POST_TELEMETRY_REQUEST", originator, msgMetaData, msgData); |
|||
} |
|||
|
|||
private void mockFindCustomer() { |
|||
when(ctxMock.getTenantId()).thenReturn(TENANT_ID); |
|||
when(ctxMock.getCustomerService()).thenReturn(customerServiceMock); |
|||
when(customerServiceMock.findCustomerByIdAsync(eq(TENANT_ID), eq(customer.getId()))).thenReturn(Futures.immediateFuture(customer)); |
|||
} |
|||
|
|||
} |
|||
@ -0,0 +1,284 @@ |
|||
/** |
|||
* Copyright © 2016-2023 The Thingsboard Authors |
|||
* |
|||
* Licensed under the Apache License, Version 2.0 (the "License"); |
|||
* you may not use this file except in compliance with the License. |
|||
* You may obtain a copy of the License at |
|||
* |
|||
* http://www.apache.org/licenses/LICENSE-2.0
|
|||
* |
|||
* Unless required by applicable law or agreed to in writing, software |
|||
* distributed under the License is distributed on an "AS IS" BASIS, |
|||
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. |
|||
* See the License for the specific language governing permissions and |
|||
* limitations under the License. |
|||
*/ |
|||
package org.thingsboard.rule.engine.metadata; |
|||
|
|||
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; |
|||
import org.mockito.ArgumentCaptor; |
|||
import org.mockito.Mock; |
|||
import org.mockito.junit.jupiter.MockitoExtension; |
|||
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.util.EntityDetails; |
|||
import org.thingsboard.server.common.data.Tenant; |
|||
import org.thingsboard.server.common.data.id.DeviceId; |
|||
import org.thingsboard.server.common.data.id.TenantId; |
|||
import org.thingsboard.server.common.msg.TbMsg; |
|||
import org.thingsboard.server.common.msg.TbMsgMetaData; |
|||
import org.thingsboard.server.dao.tenant.TenantService; |
|||
|
|||
import java.util.Collections; |
|||
import java.util.List; |
|||
import java.util.UUID; |
|||
|
|||
import static org.assertj.core.api.Assertions.assertThat; |
|||
import static org.junit.jupiter.api.Assertions.assertThrows; |
|||
import static org.mockito.ArgumentMatchers.any; |
|||
import static org.mockito.ArgumentMatchers.eq; |
|||
import static org.mockito.Mockito.never; |
|||
import static org.mockito.Mockito.times; |
|||
import static org.mockito.Mockito.verify; |
|||
import static org.mockito.Mockito.when; |
|||
|
|||
@ExtendWith(MockitoExtension.class) |
|||
public class TbGetTenantDetailsNodeTest { |
|||
|
|||
private static final DeviceId DUMMY_DEVICE_ORIGINATOR = new DeviceId(UUID.randomUUID()); |
|||
@Mock |
|||
private TbContext ctxMock; |
|||
@Mock |
|||
private TenantService tenantServiceMock; |
|||
private TbGetTenantDetailsNode node; |
|||
private TbGetTenantDetailsNodeConfiguration config; |
|||
private TbNodeConfiguration nodeConfiguration; |
|||
private TbMsg msg; |
|||
private Tenant tenant; |
|||
|
|||
@BeforeEach |
|||
public void setUp() { |
|||
node = new TbGetTenantDetailsNode(); |
|||
config = new TbGetTenantDetailsNodeConfiguration().defaultConfiguration(); |
|||
nodeConfiguration = new TbNodeConfiguration(JacksonUtil.valueToTree(config)); |
|||
tenant = new Tenant(); |
|||
tenant.setId(new TenantId(UUID.randomUUID())); |
|||
tenant.setTitle("Tenant title"); |
|||
tenant.setCountry("Tenant country"); |
|||
tenant.setCity("Tenant city"); |
|||
tenant.setState("Tenant state"); |
|||
tenant.setZip("123456"); |
|||
tenant.setAddress("Tenant address 1"); |
|||
tenant.setAddress2("Tenant address 2"); |
|||
tenant.setPhone("+123456789"); |
|||
tenant.setEmail("email@tenant.com"); |
|||
tenant.setAdditionalInfo(JacksonUtil.toJsonNode("{\"someProperty\":\"someValue\",\"description\":\"Tenant description\"}")); |
|||
} |
|||
|
|||
@Test |
|||
public void givenConfigWithNullFetchTo_whenInit_thenException() { |
|||
// GIVEN
|
|||
config.setDetailsList(List.of(EntityDetails.ID)); |
|||
config.setFetchTo(null); |
|||
nodeConfiguration = new TbNodeConfiguration(JacksonUtil.valueToTree(config)); |
|||
|
|||
// WHEN
|
|||
var exception = assertThrows(TbNodeException.class, () -> node.init(ctxMock, nodeConfiguration)); |
|||
|
|||
// THEN
|
|||
assertThat(exception.getMessage()).isEqualTo("FetchTo cannot be null!"); |
|||
verify(ctxMock, never()).tellSuccess(any()); |
|||
} |
|||
|
|||
@Test |
|||
public void givenDefaultConfig_whenInit_thenOK() { |
|||
// THEN
|
|||
assertThat(config.getDetailsList()).isEqualTo(Collections.emptyList()); |
|||
assertThat(config.getFetchTo()).isEqualTo(FetchTo.DATA); |
|||
} |
|||
|
|||
@Test |
|||
public void givenCustomConfig_whenInit_thenOK() throws TbNodeException { |
|||
// GIVEN
|
|||
config.setDetailsList(List.of(EntityDetails.ID, EntityDetails.PHONE)); |
|||
config.setFetchTo(FetchTo.METADATA); |
|||
nodeConfiguration = new TbNodeConfiguration(JacksonUtil.valueToTree(config)); |
|||
|
|||
// WHEN
|
|||
node.init(ctxMock, nodeConfiguration); |
|||
|
|||
// THEN
|
|||
assertThat(node.config).isEqualTo(config); |
|||
assertThat(config.getDetailsList()).isEqualTo(List.of(EntityDetails.ID, EntityDetails.PHONE)); |
|||
assertThat(config.getFetchTo()).isEqualTo(FetchTo.METADATA); |
|||
assertThat(node.fetchTo).isEqualTo(FetchTo.METADATA); |
|||
} |
|||
|
|||
@Test |
|||
public void givenMsgDataIsNotAnJsonObjectAndFetchToData_whenOnMsg_thenException() { |
|||
// GIVEN
|
|||
node.fetchTo = FetchTo.DATA; |
|||
msg = TbMsg.newMsg("SOME_MESSAGE_TYPE", DUMMY_DEVICE_ORIGINATOR, new TbMsgMetaData(), "[]"); |
|||
|
|||
// WHEN
|
|||
var exception = assertThrows(IllegalArgumentException.class, () -> node.onMsg(ctxMock, msg)); |
|||
|
|||
// THEN
|
|||
assertThat(exception.getMessage()).isEqualTo("Message body is not an object!"); |
|||
verify(ctxMock, never()).tellSuccess(any()); |
|||
} |
|||
|
|||
@Test |
|||
public void givenAllEntityDetailsAndFetchToData_whenOnMsg_thenShouldTellSuccessAndFetchAllToData() { |
|||
// GIVEN
|
|||
prepareMsgAndConfig(FetchTo.DATA, List.of(EntityDetails.values())); |
|||
|
|||
mockFindTenant(); |
|||
|
|||
// 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 = "{\"dataKey1\":123,\"dataKey2\":\"dataValue2\"," + |
|||
"\"tenant_id\":\"" + tenant.getId() + "\"," + |
|||
"\"tenant_title\":\"" + tenant.getTitle() + "\"," + |
|||
"\"tenant_country\":\"" + tenant.getCountry() + "\"," + |
|||
"\"tenant_city\":\"" + tenant.getCity() + "\"," + |
|||
"\"tenant_state\":\"" + tenant.getState() + "\"," + |
|||
"\"tenant_zip\":\"" + tenant.getZip() + "\"," + |
|||
"\"tenant_address\":\"" + tenant.getAddress() + "\"," + |
|||
"\"tenant_address2\":\"" + tenant.getAddress2() + "\"," + |
|||
"\"tenant_phone\":\"" + tenant.getPhone() + "\"," + |
|||
"\"tenant_email\":\"" + tenant.getEmail() + "\"," + |
|||
"\"tenant_additionalInfo\":\"" + tenant.getAdditionalInfo().get("description").asText() + "\"}"; |
|||
|
|||
assertThat(actualMessageCaptor.getValue().getData()).isEqualTo(expectedMsgData); |
|||
assertThat(actualMessageCaptor.getValue().getMetaData()).isEqualTo(msg.getMetaData()); |
|||
} |
|||
|
|||
@Test |
|||
public void givenSomeEntityDetailsAndFetchToMetadata_whenOnMsg_thenShouldTellSuccessAndFetchSomeToMetaData() { |
|||
// GIVEN
|
|||
prepareMsgAndConfig(FetchTo.METADATA, List.of(EntityDetails.ID, EntityDetails.TITLE, EntityDetails.PHONE)); |
|||
|
|||
mockFindTenant(); |
|||
|
|||
// 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(msg.getMetaData().getData()); |
|||
expectedMsgMetaData.putValue("tenant_id", tenant.getId().getId().toString()); |
|||
expectedMsgMetaData.putValue("tenant_title", tenant.getTitle()); |
|||
expectedMsgMetaData.putValue("tenant_phone", tenant.getPhone()); |
|||
|
|||
assertThat(actualMessageCaptor.getValue().getData()).isEqualTo(msg.getData()); |
|||
assertThat(actualMessageCaptor.getValue().getMetaData()).isEqualTo(expectedMsgMetaData); |
|||
} |
|||
|
|||
@Test |
|||
public void givenNotPresentEntityDetailsAndFetchToData_whenOnMsg_thenShouldTellSuccessAndFetchNothingToData() { |
|||
// GIVEN
|
|||
tenant.setZip(null); |
|||
tenant.setAddress(null); |
|||
tenant.setAddress2(null); |
|||
|
|||
prepareMsgAndConfig(FetchTo.DATA, List.of(EntityDetails.ZIP, EntityDetails.ADDRESS, EntityDetails.ADDRESS2)); |
|||
|
|||
mockFindTenant(); |
|||
|
|||
// 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()); |
|||
|
|||
assertThat(actualMessageCaptor.getValue().getData()).isEqualTo(msg.getData()); |
|||
assertThat(actualMessageCaptor.getValue().getMetaData()).isEqualTo(msg.getMetaData()); |
|||
} |
|||
|
|||
@Test |
|||
public void givenDidNotFindTenant_whenOnMsg_thenShouldTellSuccessAndFetchNothingToData() { |
|||
// GIVEN
|
|||
prepareMsgAndConfig(FetchTo.DATA, List.of(EntityDetails.ZIP, EntityDetails.ADDRESS, EntityDetails.ADDRESS2)); |
|||
|
|||
when(ctxMock.getTenantId()).thenReturn(tenant.getId()); |
|||
when(ctxMock.getTenantService()).thenReturn(tenantServiceMock); |
|||
when(tenantServiceMock.findTenantByIdAsync(eq(tenant.getId()), eq(tenant.getId()))).thenReturn(Futures.immediateFuture(null)); |
|||
|
|||
// 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()); |
|||
|
|||
assertThat(actualMessageCaptor.getValue().getData()).isEqualTo(msg.getData()); |
|||
assertThat(actualMessageCaptor.getValue().getMetaData()).isEqualTo(msg.getMetaData()); |
|||
} |
|||
|
|||
@Test |
|||
public void givenNullDescriptionAndAddInfoEntityDetails_whenOnMsg_thenShouldTellSuccessAndFetchNothingToData() { |
|||
// GIVEN
|
|||
tenant.setAdditionalInfo(JacksonUtil.toJsonNode("{\"someProperty\":\"someValue\",\"description\":null}")); |
|||
|
|||
prepareMsgAndConfig(FetchTo.DATA, List.of(EntityDetails.ADDITIONAL_INFO)); |
|||
|
|||
mockFindTenant(); |
|||
|
|||
// 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()); |
|||
|
|||
assertThat(actualMessageCaptor.getValue().getData()).isEqualTo(msg.getData()); |
|||
assertThat(actualMessageCaptor.getValue().getMetaData()).isEqualTo(msg.getMetaData()); |
|||
} |
|||
|
|||
private void prepareMsgAndConfig(FetchTo fetchTo, List<EntityDetails> detailsList) { |
|||
config.setDetailsList(detailsList); |
|||
config.setFetchTo(fetchTo); |
|||
|
|||
node.config = config; |
|||
node.fetchTo = fetchTo; |
|||
|
|||
var msgMetaData = new TbMsgMetaData(); |
|||
msgMetaData.putValue("metaKey1", "metaValue1"); |
|||
msgMetaData.putValue("metaKey2", "metaValue2"); |
|||
|
|||
var msgData = "{\"dataKey1\":123,\"dataKey2\":\"dataValue2\"}"; |
|||
|
|||
msg = TbMsg.newMsg("POST_TELEMETRY_REQUEST", DUMMY_DEVICE_ORIGINATOR, msgMetaData, msgData); |
|||
} |
|||
|
|||
private void mockFindTenant() { |
|||
when(ctxMock.getTenantId()).thenReturn(tenant.getId()); |
|||
when(ctxMock.getTenantService()).thenReturn(tenantServiceMock); |
|||
when(tenantServiceMock.findTenantByIdAsync(eq(tenant.getId()), eq(tenant.getId()))).thenReturn(Futures.immediateFuture(tenant)); |
|||
} |
|||
|
|||
} |
|||
@ -0,0 +1,170 @@ |
|||
/** |
|||
* Copyright © 2016-2023 The Thingsboard Authors |
|||
* |
|||
* Licensed under the Apache License, Version 2.0 (the "License"); |
|||
* you may not use this file except in compliance with the License. |
|||
* You may obtain a copy of the License at |
|||
* |
|||
* http://www.apache.org/licenses/LICENSE-2.0
|
|||
* |
|||
* Unless required by applicable law or agreed to in writing, software |
|||
* distributed under the License is distributed on an "AS IS" BASIS, |
|||
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. |
|||
* See the License for the specific language governing permissions and |
|||
* limitations under the License. |
|||
*/ |
|||
package org.thingsboard.rule.engine.util; |
|||
|
|||
import com.google.common.util.concurrent.Futures; |
|||
import com.google.common.util.concurrent.ListenableFuture; |
|||
import org.jetbrains.annotations.NotNull; |
|||
import org.junit.jupiter.api.Test; |
|||
import org.junit.jupiter.api.extension.ExtendWith; |
|||
import org.mockito.Mock; |
|||
import org.mockito.junit.jupiter.MockitoExtension; |
|||
import org.thingsboard.common.util.ListeningExecutor; |
|||
import org.thingsboard.rule.engine.api.TbContext; |
|||
import org.thingsboard.rule.engine.api.TbNodeException; |
|||
import org.thingsboard.server.common.data.Customer; |
|||
import org.thingsboard.server.common.data.Device; |
|||
import org.thingsboard.server.common.data.EntityType; |
|||
import org.thingsboard.server.common.data.User; |
|||
import org.thingsboard.server.common.data.asset.Asset; |
|||
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.EntityIdFactory; |
|||
import org.thingsboard.server.common.data.id.UserId; |
|||
import org.thingsboard.server.dao.asset.AssetService; |
|||
import org.thingsboard.server.dao.device.DeviceService; |
|||
import org.thingsboard.server.dao.user.UserService; |
|||
|
|||
import java.util.EnumSet; |
|||
import java.util.UUID; |
|||
import java.util.concurrent.Callable; |
|||
import java.util.concurrent.ExecutionException; |
|||
|
|||
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.mockito.ArgumentMatchers.any; |
|||
import static org.mockito.Mockito.doReturn; |
|||
import static org.mockito.Mockito.when; |
|||
|
|||
@ExtendWith(MockitoExtension.class) |
|||
public class EntitiesCustomerIdAsyncLoaderTest { |
|||
|
|||
private static final EnumSet<EntityType> SUPPORTED_ENTITY_TYPES = EnumSet.of( |
|||
EntityType.CUSTOMER, |
|||
EntityType.USER, |
|||
EntityType.ASSET, |
|||
EntityType.DEVICE |
|||
); |
|||
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 |
|||
private TbContext ctxMock; |
|||
@Mock |
|||
private UserService userServiceMock; |
|||
@Mock |
|||
private AssetService assetServiceMock; |
|||
@Mock |
|||
private DeviceService deviceServiceMock; |
|||
|
|||
@Test |
|||
public void givenCustomerEntityType_whenFindEntityIdAsync_thenOK() throws ExecutionException, InterruptedException { |
|||
// GIVEN
|
|||
var customer = new Customer(new CustomerId(UUID.randomUUID())); |
|||
|
|||
// WHEN
|
|||
var actualCustomerId = EntitiesCustomerIdAsyncLoader.findEntityIdAsync(ctxMock, customer.getId()).get(); |
|||
|
|||
// THEN
|
|||
assertEquals(customer.getId(), actualCustomerId); |
|||
} |
|||
|
|||
@Test |
|||
public void givenUserEntityType_whenFindEntityIdAsync_thenOK() throws ExecutionException, InterruptedException { |
|||
// GIVEN
|
|||
var user = new User(new UserId(UUID.randomUUID())); |
|||
var expectedCustomerId = new CustomerId(UUID.randomUUID()); |
|||
user.setCustomerId(expectedCustomerId); |
|||
|
|||
when(ctxMock.getUserService()).thenReturn(userServiceMock); |
|||
doReturn(Futures.immediateFuture(user)).when(userServiceMock).findUserByIdAsync(any(), any()); |
|||
when(ctxMock.getDbCallbackExecutor()).thenReturn(DB_EXECUTOR); |
|||
|
|||
// WHEN
|
|||
var actualCustomerId = EntitiesCustomerIdAsyncLoader.findEntityIdAsync(ctxMock, user.getId()).get(); |
|||
|
|||
// THEN
|
|||
assertEquals(expectedCustomerId, actualCustomerId); |
|||
} |
|||
|
|||
@Test |
|||
public void givenAssetEntityType_whenFindEntityIdAsync_thenOK() throws ExecutionException, InterruptedException { |
|||
// GIVEN
|
|||
var asset = new Asset(new AssetId(UUID.randomUUID())); |
|||
var expectedCustomerId = new CustomerId(UUID.randomUUID()); |
|||
asset.setCustomerId(expectedCustomerId); |
|||
|
|||
when(ctxMock.getAssetService()).thenReturn(assetServiceMock); |
|||
doReturn(Futures.immediateFuture(asset)).when(assetServiceMock).findAssetByIdAsync(any(), any()); |
|||
when(ctxMock.getDbCallbackExecutor()).thenReturn(DB_EXECUTOR); |
|||
|
|||
// WHEN
|
|||
var actualCustomerId = EntitiesCustomerIdAsyncLoader.findEntityIdAsync(ctxMock, asset.getId()).get(); |
|||
|
|||
// THEN
|
|||
assertEquals(expectedCustomerId, actualCustomerId); |
|||
} |
|||
|
|||
@Test |
|||
public void givenDeviceEntityType_whenFindEntityIdAsync_thenOK() throws ExecutionException, InterruptedException { |
|||
// GIVEN
|
|||
var device = new Device(new DeviceId(UUID.randomUUID())); |
|||
var expectedCustomerId = new CustomerId(UUID.randomUUID()); |
|||
device.setCustomerId(expectedCustomerId); |
|||
|
|||
when(ctxMock.getDeviceService()).thenReturn(deviceServiceMock); |
|||
doReturn(Futures.immediateFuture(device)).when(deviceServiceMock).findDeviceByIdAsync(any(), any()); |
|||
when(ctxMock.getDbCallbackExecutor()).thenReturn(DB_EXECUTOR); |
|||
|
|||
// WHEN
|
|||
var actualCustomerId = EntitiesCustomerIdAsyncLoader.findEntityIdAsync(ctxMock, device.getId()).get(); |
|||
|
|||
// THEN
|
|||
assertEquals(expectedCustomerId, actualCustomerId); |
|||
} |
|||
|
|||
@Test |
|||
public void givenUnsupportedEntityTypes_whenFindEntityIdAsync_thenException() { |
|||
for (var entityType : EntityType.values()) { |
|||
if (!SUPPORTED_ENTITY_TYPES.contains(entityType)) { |
|||
var entityId = EntityIdFactory.getByTypeAndUuid(entityType, UUID.randomUUID()); |
|||
|
|||
var expectedExceptionMsg = "org.thingsboard.rule.engine.api.TbNodeException: Unexpected originator EntityType: " + entityType; |
|||
|
|||
var exception = assertThrows(ExecutionException.class, |
|||
() -> EntitiesCustomerIdAsyncLoader.findEntityIdAsync(ctxMock, entityId).get()); |
|||
|
|||
assertInstanceOf(TbNodeException.class, exception.getCause()); |
|||
assertEquals(expectedExceptionMsg, exception.getMessage()); |
|||
} |
|||
} |
|||
} |
|||
|
|||
} |
|||
@ -0,0 +1,153 @@ |
|||
/** |
|||
* Copyright © 2016-2023 The Thingsboard Authors |
|||
* |
|||
* Licensed under the Apache License, Version 2.0 (the "License"); |
|||
* you may not use this file except in compliance with the License. |
|||
* You may obtain a copy of the License at |
|||
* |
|||
* http://www.apache.org/licenses/LICENSE-2.0
|
|||
* |
|||
* Unless required by applicable law or agreed to in writing, software |
|||
* distributed under the License is distributed on an "AS IS" BASIS, |
|||
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. |
|||
* See the License for the specific language governing permissions and |
|||
* limitations under the License. |
|||
*/ |
|||
package org.thingsboard.rule.engine.util; |
|||
|
|||
import com.google.common.util.concurrent.Futures; |
|||
import com.google.common.util.concurrent.ListenableFuture; |
|||
import org.jetbrains.annotations.NotNull; |
|||
import org.junit.jupiter.api.Test; |
|||
import org.junit.jupiter.api.extension.ExtendWith; |
|||
import org.mockito.Mock; |
|||
import org.mockito.junit.jupiter.MockitoExtension; |
|||
import org.thingsboard.common.util.ListeningExecutor; |
|||
import org.thingsboard.rule.engine.api.TbContext; |
|||
import org.thingsboard.rule.engine.data.DeviceRelationsQuery; |
|||
import org.thingsboard.server.common.data.Device; |
|||
import org.thingsboard.server.common.data.device.DeviceSearchQuery; |
|||
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.relation.EntityRelation; |
|||
import org.thingsboard.server.common.data.relation.EntitySearchDirection; |
|||
import org.thingsboard.server.common.data.relation.RelationsSearchParameters; |
|||
import org.thingsboard.server.dao.device.DeviceService; |
|||
|
|||
import java.util.List; |
|||
import java.util.UUID; |
|||
import java.util.concurrent.Callable; |
|||
|
|||
import static org.junit.jupiter.api.Assertions.assertEquals; |
|||
import static org.junit.jupiter.api.Assertions.assertNotNull; |
|||
import static org.mockito.ArgumentMatchers.eq; |
|||
import static org.mockito.Mockito.times; |
|||
import static org.mockito.Mockito.verify; |
|||
import static org.mockito.Mockito.when; |
|||
|
|||
@ExtendWith(MockitoExtension.class) |
|||
public class EntitiesRelatedDeviceIdAsyncLoaderTest { |
|||
|
|||
private static final EntityId DUMMY_ORIGINATOR = new DeviceId(UUID.randomUUID()); |
|||
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 |
|||
private TbContext ctxMock; |
|||
@Mock |
|||
private DeviceService deviceServiceMock; |
|||
|
|||
@Test |
|||
public void givenDeviceRelationsQuery_whenFindDeviceAsync_ShouldBuildCorrectDeviceSearchQuery() { |
|||
// GIVEN
|
|||
var deviceRelationsQuery = new DeviceRelationsQuery(); |
|||
deviceRelationsQuery.setDeviceTypes(List.of("Device type 1", "Device type 2", "default")); |
|||
deviceRelationsQuery.setDirection(EntitySearchDirection.FROM); |
|||
deviceRelationsQuery.setMaxLevel(2); |
|||
deviceRelationsQuery.setRelationType(EntityRelation.CONTAINS_TYPE); |
|||
|
|||
var expectedDeviceSearchQuery = new DeviceSearchQuery(); |
|||
var parameters = new RelationsSearchParameters( |
|||
DUMMY_ORIGINATOR, |
|||
deviceRelationsQuery.getDirection(), |
|||
deviceRelationsQuery.getMaxLevel(), |
|||
deviceRelationsQuery.isFetchLastLevelOnly() |
|||
); |
|||
expectedDeviceSearchQuery.setParameters(parameters); |
|||
expectedDeviceSearchQuery.setRelationType(deviceRelationsQuery.getRelationType()); |
|||
expectedDeviceSearchQuery.setDeviceTypes(deviceRelationsQuery.getDeviceTypes()); |
|||
|
|||
when(ctxMock.getTenantId()).thenReturn(TENANT_ID); |
|||
when(ctxMock.getDeviceService()).thenReturn(deviceServiceMock); |
|||
when(deviceServiceMock.findDevicesByQuery(eq(TENANT_ID), eq(expectedDeviceSearchQuery))) |
|||
.thenReturn(Futures.immediateFuture(null)); |
|||
when(ctxMock.getDbCallbackExecutor()).thenReturn(DB_EXECUTOR); |
|||
|
|||
// WHEN
|
|||
EntitiesRelatedDeviceIdAsyncLoader.findDeviceAsync(ctxMock, DUMMY_ORIGINATOR, deviceRelationsQuery); |
|||
|
|||
// THEN
|
|||
verify(deviceServiceMock, times(1)).findDevicesByQuery(eq(TENANT_ID), eq(expectedDeviceSearchQuery)); |
|||
} |
|||
|
|||
@Test |
|||
public void givenSeveralDevicesFound_whenFindDeviceAsync_ShouldKeepOneAndDiscardOthers() throws Exception { |
|||
// GIVEN
|
|||
var deviceRelationsQuery = new DeviceRelationsQuery(); |
|||
deviceRelationsQuery.setDeviceTypes(List.of("Device type 1", "Device type 2", "default")); |
|||
deviceRelationsQuery.setDirection(EntitySearchDirection.FROM); |
|||
deviceRelationsQuery.setMaxLevel(2); |
|||
deviceRelationsQuery.setRelationType(EntityRelation.CONTAINS_TYPE); |
|||
|
|||
var expectedDeviceSearchQuery = new DeviceSearchQuery(); |
|||
var parameters = new RelationsSearchParameters( |
|||
DUMMY_ORIGINATOR, |
|||
deviceRelationsQuery.getDirection(), |
|||
deviceRelationsQuery.getMaxLevel(), |
|||
deviceRelationsQuery.isFetchLastLevelOnly() |
|||
); |
|||
expectedDeviceSearchQuery.setParameters(parameters); |
|||
expectedDeviceSearchQuery.setRelationType(deviceRelationsQuery.getRelationType()); |
|||
expectedDeviceSearchQuery.setDeviceTypes(deviceRelationsQuery.getDeviceTypes()); |
|||
|
|||
var device1 = new Device(new DeviceId(UUID.randomUUID())); |
|||
device1.setName("Device 1"); |
|||
var device2 = new Device(new DeviceId(UUID.randomUUID())); |
|||
device1.setName("Device 2"); |
|||
var device3 = new Device(new DeviceId(UUID.randomUUID())); |
|||
device1.setName("Device 3"); |
|||
|
|||
var devicesList = List.of(device1, device2, device3); |
|||
|
|||
when(ctxMock.getTenantId()).thenReturn(TENANT_ID); |
|||
when(ctxMock.getDeviceService()).thenReturn(deviceServiceMock); |
|||
when(deviceServiceMock.findDevicesByQuery(eq(TENANT_ID), eq(expectedDeviceSearchQuery))) |
|||
.thenReturn(Futures.immediateFuture(devicesList)); |
|||
when(ctxMock.getDbCallbackExecutor()).thenReturn(DB_EXECUTOR); |
|||
|
|||
// WHEN
|
|||
var entityIdFuture = EntitiesRelatedDeviceIdAsyncLoader.findDeviceAsync(ctxMock, DUMMY_ORIGINATOR, deviceRelationsQuery); |
|||
|
|||
// THEN
|
|||
assertNotNull(entityIdFuture); |
|||
|
|||
var actualEntityId = entityIdFuture.get(); |
|||
assertNotNull(actualEntityId); |
|||
assertEquals(device1.getId(), actualEntityId); |
|||
} |
|||
|
|||
} |
|||
@ -0,0 +1,174 @@ |
|||
/** |
|||
* Copyright © 2016-2023 The Thingsboard Authors |
|||
* |
|||
* Licensed under the Apache License, Version 2.0 (the "License"); |
|||
* you may not use this file except in compliance with the License. |
|||
* You may obtain a copy of the License at |
|||
* |
|||
* http://www.apache.org/licenses/LICENSE-2.0
|
|||
* |
|||
* Unless required by applicable law or agreed to in writing, software |
|||
* distributed under the License is distributed on an "AS IS" BASIS, |
|||
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. |
|||
* See the License for the specific language governing permissions and |
|||
* limitations under the License. |
|||
*/ |
|||
package org.thingsboard.rule.engine.util; |
|||
|
|||
import com.google.common.util.concurrent.Futures; |
|||
import com.google.common.util.concurrent.ListenableFuture; |
|||
import org.jetbrains.annotations.NotNull; |
|||
import org.junit.jupiter.api.Test; |
|||
import org.junit.jupiter.api.extension.ExtendWith; |
|||
import org.mockito.Mock; |
|||
import org.mockito.junit.jupiter.MockitoExtension; |
|||
import org.thingsboard.common.util.ListeningExecutor; |
|||
import org.thingsboard.rule.engine.api.TbContext; |
|||
import org.thingsboard.rule.engine.data.RelationsQuery; |
|||
import org.thingsboard.server.common.data.Device; |
|||
import org.thingsboard.server.common.data.EntityType; |
|||
import org.thingsboard.server.common.data.asset.Asset; |
|||
import org.thingsboard.server.common.data.id.AssetId; |
|||
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.relation.EntityRelation; |
|||
import org.thingsboard.server.common.data.relation.EntityRelationsQuery; |
|||
import org.thingsboard.server.common.data.relation.EntitySearchDirection; |
|||
import org.thingsboard.server.common.data.relation.RelationEntityTypeFilter; |
|||
import org.thingsboard.server.common.data.relation.RelationsSearchParameters; |
|||
import org.thingsboard.server.dao.relation.RelationService; |
|||
|
|||
import java.util.Collections; |
|||
import java.util.List; |
|||
import java.util.UUID; |
|||
import java.util.concurrent.Callable; |
|||
|
|||
import static org.junit.jupiter.api.Assertions.assertEquals; |
|||
import static org.junit.jupiter.api.Assertions.assertNotNull; |
|||
import static org.mockito.ArgumentMatchers.eq; |
|||
import static org.mockito.Mockito.times; |
|||
import static org.mockito.Mockito.verify; |
|||
import static org.mockito.Mockito.when; |
|||
|
|||
@ExtendWith(MockitoExtension.class) |
|||
public class EntitiesRelatedEntitiesIdAsyncLoaderTest { |
|||
|
|||
private static final EntityId DUMMY_ORIGINATOR = new DeviceId(UUID.randomUUID()); |
|||
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 |
|||
private TbContext ctxMock; |
|||
@Mock |
|||
private RelationService relationServiceMock; |
|||
|
|||
@Test |
|||
public void givenRelationsQuery_whenFindEntityAsync_ShouldBuildCorrectEntityRelationsQuery() { |
|||
// GIVEN
|
|||
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)); |
|||
|
|||
var expectedEntityRelationsQuery = new EntityRelationsQuery(); |
|||
var parameters = new RelationsSearchParameters( |
|||
DUMMY_ORIGINATOR, |
|||
relationsQuery.getDirection(), |
|||
relationsQuery.getMaxLevel(), |
|||
relationsQuery.isFetchLastLevelOnly() |
|||
); |
|||
expectedEntityRelationsQuery.setParameters(parameters); |
|||
expectedEntityRelationsQuery.setFilters(relationsQuery.getFilters()); |
|||
|
|||
when(ctxMock.getTenantId()).thenReturn(TENANT_ID); |
|||
when(ctxMock.getRelationService()).thenReturn(relationServiceMock); |
|||
when(relationServiceMock.findByQuery(eq(TENANT_ID), eq(expectedEntityRelationsQuery))) |
|||
.thenReturn(Futures.immediateFuture(null)); |
|||
when(ctxMock.getDbCallbackExecutor()).thenReturn(DB_EXECUTOR); |
|||
|
|||
// WHEN
|
|||
EntitiesRelatedEntityIdAsyncLoader.findEntityAsync(ctxMock, DUMMY_ORIGINATOR, relationsQuery); |
|||
|
|||
// THEN
|
|||
verify(relationServiceMock, times(1)).findByQuery(eq(TENANT_ID), eq(expectedEntityRelationsQuery)); |
|||
} |
|||
|
|||
@Test |
|||
public void givenSeveralEntitiesFound_whenFindEntityAsync_ShouldKeepOneAndDiscardOthers() throws Exception { |
|||
// GIVEN
|
|||
var relationsQuery = new RelationsQuery(); |
|||
var relationEntityTypeFilter = new RelationEntityTypeFilter( |
|||
EntityRelation.CONTAINS_TYPE, |
|||
List.of(EntityType.DEVICE, EntityType.ASSET) |
|||
); |
|||
relationsQuery.setDirection(EntitySearchDirection.FROM); |
|||
relationsQuery.setMaxLevel(2); |
|||
relationsQuery.setFilters(Collections.singletonList(relationEntityTypeFilter)); |
|||
|
|||
var expectedEntityRelationsQuery = new EntityRelationsQuery(); |
|||
var parameters = new RelationsSearchParameters( |
|||
DUMMY_ORIGINATOR, |
|||
relationsQuery.getDirection(), |
|||
relationsQuery.getMaxLevel(), |
|||
relationsQuery.isFetchLastLevelOnly() |
|||
); |
|||
expectedEntityRelationsQuery.setParameters(parameters); |
|||
expectedEntityRelationsQuery.setFilters(relationsQuery.getFilters()); |
|||
|
|||
var device1 = new Device(new DeviceId(UUID.randomUUID())); |
|||
device1.setName("Device 1"); |
|||
var device2 = new Device(new DeviceId(UUID.randomUUID())); |
|||
device1.setName("Device 2"); |
|||
var asset = new Asset(new AssetId(UUID.randomUUID())); |
|||
asset.setName("Asset"); |
|||
|
|||
var entityRelationDevice1 = new EntityRelation(); |
|||
entityRelationDevice1.setFrom(DUMMY_ORIGINATOR); |
|||
entityRelationDevice1.setTo(device1.getId()); |
|||
entityRelationDevice1.setType(EntityRelation.CONTAINS_TYPE); |
|||
|
|||
var entityRelationDevice2 = new EntityRelation(); |
|||
entityRelationDevice2.setFrom(DUMMY_ORIGINATOR); |
|||
entityRelationDevice2.setTo(device2.getId()); |
|||
entityRelationDevice2.setType(EntityRelation.CONTAINS_TYPE); |
|||
|
|||
var entityRelationAsset = new EntityRelation(); |
|||
entityRelationAsset.setFrom(DUMMY_ORIGINATOR); |
|||
entityRelationAsset.setTo(asset.getId()); |
|||
entityRelationAsset.setType(EntityRelation.CONTAINS_TYPE); |
|||
|
|||
var expectedEntityRelationsList = List.of(entityRelationDevice1, entityRelationDevice2, entityRelationAsset); |
|||
|
|||
when(ctxMock.getTenantId()).thenReturn(TENANT_ID); |
|||
when(ctxMock.getRelationService()).thenReturn(relationServiceMock); |
|||
when(relationServiceMock.findByQuery(eq(TENANT_ID), eq(expectedEntityRelationsQuery))) |
|||
.thenReturn(Futures.immediateFuture(expectedEntityRelationsList)); |
|||
when(ctxMock.getDbCallbackExecutor()).thenReturn(DB_EXECUTOR); |
|||
|
|||
// WHEN
|
|||
var deviceIdFuture = EntitiesRelatedEntityIdAsyncLoader.findEntityAsync(ctxMock, DUMMY_ORIGINATOR, relationsQuery); |
|||
|
|||
// THEN
|
|||
assertNotNull(deviceIdFuture); |
|||
|
|||
var actualDeviceId = deviceIdFuture.get(); |
|||
assertNotNull(actualDeviceId); |
|||
assertEquals(device1.getId(), actualDeviceId); |
|||
} |
|||
|
|||
} |
|||
Loading…
Reference in new issue