|
|
|
@ -16,7 +16,7 @@ |
|
|
|
package org.thingsboard.rule.engine.metadata; |
|
|
|
|
|
|
|
import com.fasterxml.jackson.databind.JsonNode; |
|
|
|
import com.google.common.util.concurrent.Futures; |
|
|
|
import com.google.common.util.concurrent.FluentFuture; |
|
|
|
import lombok.RequiredArgsConstructor; |
|
|
|
import org.junit.jupiter.api.Assertions; |
|
|
|
import org.junit.jupiter.api.BeforeEach; |
|
|
|
@ -53,19 +53,19 @@ import org.thingsboard.server.common.data.msg.TbMsgType; |
|
|
|
import org.thingsboard.server.common.data.util.TbPair; |
|
|
|
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.attributes.AttributesService; |
|
|
|
import org.thingsboard.server.dao.device.DeviceService; |
|
|
|
import org.thingsboard.server.dao.entity.EntityService; |
|
|
|
import org.thingsboard.server.dao.timeseries.TimeseriesService; |
|
|
|
import org.thingsboard.server.dao.user.UserService; |
|
|
|
|
|
|
|
import java.util.Arrays; |
|
|
|
import java.util.Collections; |
|
|
|
import java.util.List; |
|
|
|
import java.util.Map; |
|
|
|
import java.util.NoSuchElementException; |
|
|
|
import java.util.Optional; |
|
|
|
import java.util.UUID; |
|
|
|
|
|
|
|
import static com.google.common.util.concurrent.Futures.immediateFuture; |
|
|
|
import static org.assertj.core.api.Assertions.assertThat; |
|
|
|
import static org.junit.jupiter.api.Assertions.assertEquals; |
|
|
|
import static org.junit.jupiter.api.Assertions.assertInstanceOf; |
|
|
|
@ -73,19 +73,20 @@ import static org.junit.jupiter.api.Assertions.assertThrows; |
|
|
|
import static org.mockito.ArgumentMatchers.any; |
|
|
|
import static org.mockito.ArgumentMatchers.argThat; |
|
|
|
import static org.mockito.ArgumentMatchers.eq; |
|
|
|
import static org.mockito.Mockito.doReturn; |
|
|
|
import static org.mockito.Mockito.lenient; |
|
|
|
import static org.mockito.Mockito.never; |
|
|
|
import static org.mockito.Mockito.times; |
|
|
|
import static org.mockito.Mockito.verify; |
|
|
|
import static org.mockito.Mockito.verifyNoInteractions; |
|
|
|
import static org.mockito.Mockito.when; |
|
|
|
|
|
|
|
@ExtendWith(MockitoExtension.class) |
|
|
|
public class TbGetCustomerAttributeNodeTest { |
|
|
|
|
|
|
|
private static final DeviceId DUMMY_DEVICE_ORIGINATOR = new DeviceId(UUID.randomUUID()); |
|
|
|
private static final TenantId TENANT_ID = new TenantId(UUID.randomUUID()); |
|
|
|
private static final CustomerId CUSTOMER_ID = new CustomerId(UUID.randomUUID()); |
|
|
|
private static final ListeningExecutor DB_EXECUTOR = new TestDbCallbackExecutor(); |
|
|
|
private final DeviceId DUMMY_DEVICE_ORIGINATOR = new DeviceId(UUID.randomUUID()); |
|
|
|
private final TenantId TENANT_ID = TenantId.fromUUID(UUID.randomUUID()); |
|
|
|
private final CustomerId CUSTOMER_ID = new CustomerId(UUID.randomUUID()); |
|
|
|
private final ListeningExecutor DB_EXECUTOR = new TestDbCallbackExecutor(); |
|
|
|
|
|
|
|
@Mock |
|
|
|
private TbContext ctxMock; |
|
|
|
@Mock |
|
|
|
@ -93,21 +94,24 @@ public class TbGetCustomerAttributeNodeTest { |
|
|
|
@Mock |
|
|
|
private TimeseriesService timeseriesServiceMock; |
|
|
|
@Mock |
|
|
|
private UserService userServiceMock; |
|
|
|
@Mock |
|
|
|
private AssetService assetServiceMock; |
|
|
|
@Mock |
|
|
|
private DeviceService deviceServiceMock; |
|
|
|
private EntityService entityServiceMock; |
|
|
|
|
|
|
|
private TbGetCustomerAttributeNode node; |
|
|
|
private TbGetEntityDataNodeConfiguration config; |
|
|
|
private TbNodeConfiguration nodeConfiguration; |
|
|
|
private TbMsg msg; |
|
|
|
|
|
|
|
@BeforeEach |
|
|
|
public void setUp() { |
|
|
|
public void setup() { |
|
|
|
node = new TbGetCustomerAttributeNode(); |
|
|
|
config = new TbGetEntityDataNodeConfiguration().defaultConfiguration(); |
|
|
|
nodeConfiguration = new TbNodeConfiguration(JacksonUtil.valueToTree(config)); |
|
|
|
|
|
|
|
lenient().when(ctxMock.getTenantId()).thenReturn(TENANT_ID); |
|
|
|
lenient().when(ctxMock.getDbCallbackExecutor()).thenReturn(DB_EXECUTOR); |
|
|
|
lenient().when(ctxMock.getAttributesService()).thenReturn(attributesServiceMock); |
|
|
|
lenient().when(ctxMock.getTimeseriesService()).thenReturn(timeseriesServiceMock); |
|
|
|
lenient().when(ctxMock.getEntityService()).thenReturn(entityServiceMock); |
|
|
|
} |
|
|
|
|
|
|
|
@Test |
|
|
|
@ -237,12 +241,9 @@ public class TbGetCustomerAttributeNodeTest { |
|
|
|
.data(TbMsg.EMPTY_JSON_OBJECT) |
|
|
|
.build(); |
|
|
|
|
|
|
|
when(ctxMock.getTenantId()).thenReturn(TENANT_ID); |
|
|
|
|
|
|
|
when(ctxMock.getUserService()).thenReturn(userServiceMock); |
|
|
|
doReturn(Futures.immediateFuture(null)).when(userServiceMock).findUserByIdAsync(eq(TENANT_ID), eq(userId)); |
|
|
|
|
|
|
|
when(ctxMock.getDbCallbackExecutor()).thenReturn(DB_EXECUTOR); |
|
|
|
when(entityServiceMock.fetchEntityCustomerIdAsync(TENANT_ID, userId)).thenReturn( |
|
|
|
FluentFuture.from(immediateFuture(Optional.empty())) |
|
|
|
); |
|
|
|
|
|
|
|
// WHEN
|
|
|
|
node.onMsg(ctxMock, msg); |
|
|
|
@ -252,18 +253,13 @@ public class TbGetCustomerAttributeNodeTest { |
|
|
|
var actualExceptionCaptor = ArgumentCaptor.forClass(Throwable.class); |
|
|
|
|
|
|
|
verify(ctxMock, never()).tellSuccess(any()); |
|
|
|
verify(ctxMock, times(1)) |
|
|
|
.tellFailure(actualMessageCaptor.capture(), actualExceptionCaptor.capture()); |
|
|
|
verify(ctxMock).tellFailure(actualMessageCaptor.capture(), actualExceptionCaptor.capture()); |
|
|
|
|
|
|
|
var actualMessage = actualMessageCaptor.getValue(); |
|
|
|
var actualException = actualExceptionCaptor.getValue(); |
|
|
|
|
|
|
|
var expectedExceptionMessage = String.format( |
|
|
|
"Failed to find customer for entity with id: %s and type: %s", |
|
|
|
userId.getId(), userId.getEntityType().getNormalName()); |
|
|
|
|
|
|
|
assertEquals(msg, actualMessage); |
|
|
|
assertEquals(expectedExceptionMessage, actualException.getMessage()); |
|
|
|
assertEquals("Originator not found", actualException.getMessage()); |
|
|
|
assertInstanceOf(NoSuchElementException.class, actualException); |
|
|
|
} |
|
|
|
|
|
|
|
@ -282,16 +278,12 @@ public class TbGetCustomerAttributeNodeTest { |
|
|
|
); |
|
|
|
var expectedPatternProcessedKeysList = List.of("sourceKey1", "sourceKey2", "sourceKey3"); |
|
|
|
|
|
|
|
when(ctxMock.getTenantId()).thenReturn(TENANT_ID); |
|
|
|
|
|
|
|
when(ctxMock.getDeviceService()).thenReturn(deviceServiceMock); |
|
|
|
doReturn(device).when(deviceServiceMock).findDeviceById(eq(TENANT_ID), eq(device.getId())); |
|
|
|
when(entityServiceMock.fetchEntityCustomerIdAsync(TENANT_ID, device.getId())).thenReturn( |
|
|
|
FluentFuture.from(immediateFuture(Optional.of(CUSTOMER_ID))) |
|
|
|
); |
|
|
|
|
|
|
|
when(ctxMock.getAttributesService()).thenReturn(attributesServiceMock); |
|
|
|
when(attributesServiceMock.find(eq(TENANT_ID), eq(CUSTOMER_ID), eq(AttributeScope.SERVER_SCOPE), argThat(new ListMatcher<>(expectedPatternProcessedKeysList)))) |
|
|
|
.thenReturn(Futures.immediateFuture(attributesList)); |
|
|
|
|
|
|
|
when(ctxMock.getDbCallbackExecutor()).thenReturn(DB_EXECUTOR); |
|
|
|
.thenReturn(immediateFuture(attributesList)); |
|
|
|
|
|
|
|
// WHEN
|
|
|
|
node.onMsg(ctxMock, msg); |
|
|
|
@ -299,7 +291,7 @@ public class TbGetCustomerAttributeNodeTest { |
|
|
|
// THEN
|
|
|
|
var actualMessageCaptor = ArgumentCaptor.forClass(TbMsg.class); |
|
|
|
|
|
|
|
verify(ctxMock, times(1)).tellSuccess(actualMessageCaptor.capture()); |
|
|
|
verify(ctxMock).tellSuccess(actualMessageCaptor.capture()); |
|
|
|
verify(ctxMock, never()).tellFailure(any(), any()); |
|
|
|
|
|
|
|
var expectedMsgData = "{\"temp\":42," + |
|
|
|
@ -329,16 +321,12 @@ public class TbGetCustomerAttributeNodeTest { |
|
|
|
); |
|
|
|
var expectedPatternProcessedKeysList = List.of("sourceKey1", "sourceKey2", "sourceKey3"); |
|
|
|
|
|
|
|
when(ctxMock.getTenantId()).thenReturn(TENANT_ID); |
|
|
|
|
|
|
|
when(ctxMock.getUserService()).thenReturn(userServiceMock); |
|
|
|
doReturn(Futures.immediateFuture(user)).when(userServiceMock).findUserByIdAsync(eq(TENANT_ID), eq(user.getId())); |
|
|
|
when(entityServiceMock.fetchEntityCustomerIdAsync(TENANT_ID, user.getId())).thenReturn( |
|
|
|
FluentFuture.from(immediateFuture(Optional.of(CUSTOMER_ID))) |
|
|
|
); |
|
|
|
|
|
|
|
when(ctxMock.getAttributesService()).thenReturn(attributesServiceMock); |
|
|
|
when(attributesServiceMock.find(eq(TENANT_ID), eq(CUSTOMER_ID), eq(AttributeScope.SERVER_SCOPE), argThat(new ListMatcher<>(expectedPatternProcessedKeysList)))) |
|
|
|
.thenReturn(Futures.immediateFuture(attributesList)); |
|
|
|
|
|
|
|
when(ctxMock.getDbCallbackExecutor()).thenReturn(DB_EXECUTOR); |
|
|
|
.thenReturn(immediateFuture(attributesList)); |
|
|
|
|
|
|
|
// WHEN
|
|
|
|
node.onMsg(ctxMock, msg); |
|
|
|
@ -346,7 +334,7 @@ public class TbGetCustomerAttributeNodeTest { |
|
|
|
// THEN
|
|
|
|
var actualMessageCaptor = ArgumentCaptor.forClass(TbMsg.class); |
|
|
|
|
|
|
|
verify(ctxMock, times(1)).tellSuccess(actualMessageCaptor.capture()); |
|
|
|
verify(ctxMock).tellSuccess(actualMessageCaptor.capture()); |
|
|
|
verify(ctxMock, never()).tellFailure(any(), any()); |
|
|
|
|
|
|
|
var expectedMsgMetaData = new TbMsgMetaData(Map.of( |
|
|
|
@ -375,21 +363,18 @@ public class TbGetCustomerAttributeNodeTest { |
|
|
|
); |
|
|
|
var expectedPatternProcessedKeysList = List.of("sourceKey1", "sourceKey2", "sourceKey3"); |
|
|
|
|
|
|
|
when(ctxMock.getTenantId()).thenReturn(TENANT_ID); |
|
|
|
|
|
|
|
when(ctxMock.getTimeseriesService()).thenReturn(timeseriesServiceMock); |
|
|
|
when(timeseriesServiceMock.findLatest(eq(TENANT_ID), eq(customer.getId()), argThat(new ListMatcher<>(expectedPatternProcessedKeysList)))) |
|
|
|
.thenReturn(Futures.immediateFuture(timeseriesList)); |
|
|
|
|
|
|
|
when(ctxMock.getDbCallbackExecutor()).thenReturn(DB_EXECUTOR); |
|
|
|
.thenReturn(immediateFuture(timeseriesList)); |
|
|
|
|
|
|
|
// WHEN
|
|
|
|
node.onMsg(ctxMock, msg); |
|
|
|
|
|
|
|
// THEN
|
|
|
|
verifyNoInteractions(entityServiceMock); |
|
|
|
|
|
|
|
var actualMessageCaptor = ArgumentCaptor.forClass(TbMsg.class); |
|
|
|
|
|
|
|
verify(ctxMock, times(1)).tellSuccess(actualMessageCaptor.capture()); |
|
|
|
verify(ctxMock).tellSuccess(actualMessageCaptor.capture()); |
|
|
|
verify(ctxMock, never()).tellFailure(any(), any()); |
|
|
|
|
|
|
|
var expectedMsgData = "{\"temp\":42," + |
|
|
|
@ -408,7 +393,7 @@ public class TbGetCustomerAttributeNodeTest { |
|
|
|
public void givenFetchTelemetryToMetaData_whenOnMsg_thenShouldFetchTelemetryToMetaData() { |
|
|
|
// GIVEN
|
|
|
|
var asset = new Asset(new AssetId(UUID.randomUUID())); |
|
|
|
asset.setCustomerId(new CustomerId(UUID.randomUUID())); |
|
|
|
asset.setCustomerId(CUSTOMER_ID); |
|
|
|
|
|
|
|
prepareMsgAndConfig(TbMsgSource.METADATA, DataToFetch.LATEST_TELEMETRY, asset.getId()); |
|
|
|
|
|
|
|
@ -419,16 +404,12 @@ public class TbGetCustomerAttributeNodeTest { |
|
|
|
); |
|
|
|
var expectedPatternProcessedKeysList = List.of("sourceKey1", "sourceKey2", "sourceKey3"); |
|
|
|
|
|
|
|
when(ctxMock.getTenantId()).thenReturn(TENANT_ID); |
|
|
|
|
|
|
|
when(ctxMock.getAssetService()).thenReturn(assetServiceMock); |
|
|
|
doReturn(Futures.immediateFuture(asset)).when(assetServiceMock).findAssetByIdAsync(eq(TENANT_ID), eq(asset.getId())); |
|
|
|
|
|
|
|
when(ctxMock.getTimeseriesService()).thenReturn(timeseriesServiceMock); |
|
|
|
when(timeseriesServiceMock.findLatest(eq(TENANT_ID), eq(asset.getCustomerId()), argThat(new ListMatcher<>(expectedPatternProcessedKeysList)))) |
|
|
|
.thenReturn(Futures.immediateFuture(timeseriesList)); |
|
|
|
when(entityServiceMock.fetchEntityCustomerIdAsync(TENANT_ID, asset.getId())).thenReturn( |
|
|
|
FluentFuture.from(immediateFuture(Optional.of(CUSTOMER_ID))) |
|
|
|
); |
|
|
|
|
|
|
|
when(ctxMock.getDbCallbackExecutor()).thenReturn(DB_EXECUTOR); |
|
|
|
when(timeseriesServiceMock.findLatest(eq(TENANT_ID), eq(CUSTOMER_ID), argThat(new ListMatcher<>(expectedPatternProcessedKeysList)))) |
|
|
|
.thenReturn(immediateFuture(timeseriesList)); |
|
|
|
|
|
|
|
// WHEN
|
|
|
|
node.onMsg(ctxMock, msg); |
|
|
|
@ -436,7 +417,7 @@ public class TbGetCustomerAttributeNodeTest { |
|
|
|
// THEN
|
|
|
|
var actualMessageCaptor = ArgumentCaptor.forClass(TbMsg.class); |
|
|
|
|
|
|
|
verify(ctxMock, times(1)).tellSuccess(actualMessageCaptor.capture()); |
|
|
|
verify(ctxMock).tellSuccess(actualMessageCaptor.capture()); |
|
|
|
verify(ctxMock, never()).tellFailure(any(), any()); |
|
|
|
|
|
|
|
var expectedMsgMetaData = new TbMsgMetaData(Map.of( |
|
|
|
|