Browse Source

Review fixes: replace direct executor with db executor

pull/8661/head
Dmytro Skarzhynets 4 years ago
parent
commit
913c6b60e2
  1. 6
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbAbstractGetAttributesNode.java
  2. 7
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbAbstractGetEntityAttrNode.java
  3. 11
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetCustomerDetailsNode.java
  4. 11
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/util/EntitiesCustomerIdAsyncLoader.java
  5. 3
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/util/EntitiesRelatedDeviceIdAsyncLoader.java
  6. 5
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/util/EntitiesRelatedEntityIdAsyncLoader.java
  7. 38
      rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/metadata/TbGetCustomerDetailsNodeTest.java
  8. 31
      rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/util/EntitiesCustomerIdAsyncLoaderTest.java
  9. 21
      rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/util/EntitiesRelatedDeviceIdAsyncLoaderTest.java
  10. 21
      rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/util/EntitiesRelatedEntitiesIdAsyncLoaderTest.java

6
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbAbstractGetAttributesNode.java

@ -101,7 +101,7 @@ public abstract class TbAbstractGetAttributesNode<C extends TbGetAttributesNodeC
} else { } else {
ctx.tellFailure(outMsg, reportFailures(failuresPairSet)); ctx.tellFailure(outMsg, reportFailures(failuresPairSet));
} }
}, t -> ctx.tellFailure(msg, t), ctx.getDbCallbackExecutor()); }, t -> ctx.tellFailure(msg, t), MoreExecutors.directExecutor());
} }
private ListenableFuture<TbPair<String, List<AttributeKvEntry>>> getAttrAsync( private ListenableFuture<TbPair<String, List<AttributeKvEntry>>> getAttrAsync(
@ -121,7 +121,7 @@ public abstract class TbAbstractGetAttributesNode<C extends TbGetAttributesNodeC
failuresPairSet.add(new TbPair<>(scope, nonExistentKeys)); failuresPairSet.add(new TbPair<>(scope, nonExistentKeys));
} }
return new TbPair<>(scope, attributeKvEntryList); return new TbPair<>(scope, attributeKvEntryList);
}, MoreExecutors.directExecutor()); }, ctx.getDbCallbackExecutor());
} }
private ListenableFuture<TbPair<String, List<TsKvEntry>>> getLatestTelemetry(TbContext ctx, EntityId entityId, List<String> keys, Set<TbPair<String, List<String>>> failuresPairSet) { private ListenableFuture<TbPair<String, List<TsKvEntry>>> getLatestTelemetry(TbContext ctx, EntityId entityId, List<String> keys, Set<TbPair<String, List<String>>> failuresPairSet) {
@ -147,7 +147,7 @@ public abstract class TbAbstractGetAttributesNode<C extends TbGetAttributesNodeC
failuresPairSet.add(new TbPair<>(LATEST_TS, nonExistentKeys)); failuresPairSet.add(new TbPair<>(LATEST_TS, nonExistentKeys));
} }
return new TbPair<>(LATEST_TS, listTsKvEntry); return new TbPair<>(LATEST_TS, listTsKvEntry);
}, MoreExecutors.directExecutor()); }, ctx.getDbCallbackExecutor());
} }
private TsKvEntry getValueWithTs(TsKvEntry tsKvEntry) { private TsKvEntry getValueWithTs(TsKvEntry tsKvEntry) {

7
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbAbstractGetEntityAttrNode.java

@ -65,7 +65,8 @@ public abstract class TbAbstractGetEntityAttrNode<T extends EntityId> extends Tb
var sourceKeys = List.copyOf(mappingsMap.keySet()); var sourceKeys = List.copyOf(mappingsMap.keySet());
withCallback(config.isTelemetry() ? getLatestTelemetryAsync(ctx, entityId, sourceKeys) : getAttributesAsync(ctx, entityId, sourceKeys), withCallback(config.isTelemetry() ? getLatestTelemetryAsync(ctx, entityId, sourceKeys) : getAttributesAsync(ctx, entityId, sourceKeys),
data -> putDataAndTell(ctx, msg, data, mappingsMap, msgDataAsJsonNode), data -> putDataAndTell(ctx, msg, data, mappingsMap, msgDataAsJsonNode),
t -> ctx.tellFailure(msg, t), ctx.getDbCallbackExecutor()); t -> ctx.tellFailure(msg, t),
MoreExecutors.directExecutor());
} }
private ListenableFuture<List<KvEntry>> getAttributesAsync(TbContext ctx, EntityId entityId, List<String> attrKeys) { private ListenableFuture<List<KvEntry>> getAttributesAsync(TbContext ctx, EntityId entityId, List<String> attrKeys) {
@ -74,7 +75,7 @@ public abstract class TbAbstractGetEntityAttrNode<T extends EntityId> extends Tb
l.stream() l.stream()
.map(i -> (KvEntry) i) .map(i -> (KvEntry) i)
.collect(Collectors.toList()), .collect(Collectors.toList()),
MoreExecutors.directExecutor()); ctx.getDbCallbackExecutor());
} }
private ListenableFuture<List<KvEntry>> getLatestTelemetryAsync(TbContext ctx, EntityId entityId, List<String> timeseriesKeys) { private ListenableFuture<List<KvEntry>> getLatestTelemetryAsync(TbContext ctx, EntityId entityId, List<String> timeseriesKeys) {
@ -83,7 +84,7 @@ public abstract class TbAbstractGetEntityAttrNode<T extends EntityId> extends Tb
l.stream() l.stream()
.map(i -> (KvEntry) i) .map(i -> (KvEntry) i)
.collect(Collectors.toList()), .collect(Collectors.toList()),
MoreExecutors.directExecutor()); ctx.getDbCallbackExecutor());
} }
private void putDataAndTell(TbContext ctx, TbMsg msg, List<? extends KvEntry> data, Map<String, String> map, ObjectNode msgData) { private void putDataAndTell(TbContext ctx, TbMsg msg, List<? extends KvEntry> data, Map<String, String> map, ObjectNode msgData) {

11
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetCustomerDetailsNode.java

@ -17,7 +17,6 @@ package org.thingsboard.rule.engine.metadata;
import com.google.common.util.concurrent.Futures; import com.google.common.util.concurrent.Futures;
import com.google.common.util.concurrent.ListenableFuture; import com.google.common.util.concurrent.ListenableFuture;
import com.google.common.util.concurrent.MoreExecutors;
import lombok.extern.slf4j.Slf4j; import lombok.extern.slf4j.Slf4j;
import org.thingsboard.rule.engine.api.RuleNode; import org.thingsboard.rule.engine.api.RuleNode;
import org.thingsboard.rule.engine.api.TbContext; import org.thingsboard.rule.engine.api.TbContext;
@ -70,19 +69,19 @@ public class TbGetCustomerDetailsNode extends TbAbstractGetEntityDetailsNode<TbG
switch (msg.getOriginator().getEntityType()) { switch (msg.getOriginator().getEntityType()) {
case DEVICE: case DEVICE:
return Futures.transformAsync(ctx.getDeviceService().findDeviceByIdAsync(ctx.getTenantId(), new DeviceId(msg.getOriginator().getId())), return Futures.transformAsync(ctx.getDeviceService().findDeviceByIdAsync(ctx.getTenantId(), new DeviceId(msg.getOriginator().getId())),
device -> getCustomerFuture(ctx, device, msg.getOriginator()), MoreExecutors.directExecutor()); device -> getCustomerFuture(ctx, device, msg.getOriginator()), ctx.getDbCallbackExecutor());
case ASSET: case ASSET:
return Futures.transformAsync(ctx.getAssetService().findAssetByIdAsync(ctx.getTenantId(), new AssetId(msg.getOriginator().getId())), return Futures.transformAsync(ctx.getAssetService().findAssetByIdAsync(ctx.getTenantId(), new AssetId(msg.getOriginator().getId())),
asset -> getCustomerFuture(ctx, asset, msg.getOriginator()), MoreExecutors.directExecutor()); asset -> getCustomerFuture(ctx, asset, msg.getOriginator()), ctx.getDbCallbackExecutor());
case ENTITY_VIEW: case ENTITY_VIEW:
return Futures.transformAsync(ctx.getEntityViewService().findEntityViewByIdAsync(ctx.getTenantId(), new EntityViewId(msg.getOriginator().getId())), return Futures.transformAsync(ctx.getEntityViewService().findEntityViewByIdAsync(ctx.getTenantId(), new EntityViewId(msg.getOriginator().getId())),
entityView -> getCustomerFuture(ctx, entityView, msg.getOriginator()), MoreExecutors.directExecutor()); entityView -> getCustomerFuture(ctx, entityView, msg.getOriginator()), ctx.getDbCallbackExecutor());
case USER: case USER:
return Futures.transformAsync(ctx.getUserService().findUserByIdAsync(ctx.getTenantId(), new UserId(msg.getOriginator().getId())), return Futures.transformAsync(ctx.getUserService().findUserByIdAsync(ctx.getTenantId(), new UserId(msg.getOriginator().getId())),
user -> getCustomerFuture(ctx, user, msg.getOriginator()), MoreExecutors.directExecutor()); user -> getCustomerFuture(ctx, user, msg.getOriginator()), ctx.getDbCallbackExecutor());
case EDGE: case EDGE:
return Futures.transformAsync(ctx.getEdgeService().findEdgeByIdAsync(ctx.getTenantId(), new EdgeId(msg.getOriginator().getId())), return Futures.transformAsync(ctx.getEdgeService().findEdgeByIdAsync(ctx.getTenantId(), new EdgeId(msg.getOriginator().getId())),
edge -> getCustomerFuture(ctx, edge, msg.getOriginator()), MoreExecutors.directExecutor()); edge -> getCustomerFuture(ctx, edge, msg.getOriginator()), ctx.getDbCallbackExecutor());
default: default:
return Futures.immediateFailedFuture(new NoSuchElementException("Entity with entityType '" + msg.getOriginator().getEntityType() + "' is not supported.")); return Futures.immediateFailedFuture(new NoSuchElementException("Entity with entityType '" + msg.getOriginator().getEntityType() + "' is not supported."));
} }

11
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/util/EntitiesCustomerIdAsyncLoader.java

@ -17,7 +17,6 @@ package org.thingsboard.rule.engine.util;
import com.google.common.util.concurrent.Futures; import com.google.common.util.concurrent.Futures;
import com.google.common.util.concurrent.ListenableFuture; import com.google.common.util.concurrent.ListenableFuture;
import com.google.common.util.concurrent.MoreExecutors;
import org.thingsboard.rule.engine.api.TbContext; import org.thingsboard.rule.engine.api.TbContext;
import org.thingsboard.rule.engine.api.TbNodeException; import org.thingsboard.rule.engine.api.TbNodeException;
import org.thingsboard.server.common.data.HasCustomerId; import org.thingsboard.server.common.data.HasCustomerId;
@ -34,19 +33,19 @@ public class EntitiesCustomerIdAsyncLoader {
case CUSTOMER: case CUSTOMER:
return Futures.immediateFuture((CustomerId) originator); return Futures.immediateFuture((CustomerId) originator);
case USER: case USER:
return toCustomerIdAsync(ctx.getUserService().findUserByIdAsync(ctx.getTenantId(), (UserId) originator)); return toCustomerIdAsync(ctx, ctx.getUserService().findUserByIdAsync(ctx.getTenantId(), (UserId) originator));
case ASSET: case ASSET:
return toCustomerIdAsync(ctx.getAssetService().findAssetByIdAsync(ctx.getTenantId(), (AssetId) originator)); return toCustomerIdAsync(ctx, ctx.getAssetService().findAssetByIdAsync(ctx.getTenantId(), (AssetId) originator));
case DEVICE: case DEVICE:
return toCustomerIdAsync(ctx.getDeviceService().findDeviceByIdAsync(ctx.getTenantId(), (DeviceId) originator)); return toCustomerIdAsync(ctx, ctx.getDeviceService().findDeviceByIdAsync(ctx.getTenantId(), (DeviceId) originator));
default: default:
return Futures.immediateFailedFuture(new TbNodeException("Unexpected originator EntityType: " + originator.getEntityType())); return Futures.immediateFailedFuture(new TbNodeException("Unexpected originator EntityType: " + originator.getEntityType()));
} }
} }
private static <T extends HasCustomerId> ListenableFuture<CustomerId> toCustomerIdAsync(ListenableFuture<T> future) { private static <T extends HasCustomerId> ListenableFuture<CustomerId> toCustomerIdAsync(TbContext ctx, ListenableFuture<T> future) {
return Futures.transformAsync(future, in -> in != null ? Futures.immediateFuture(in.getCustomerId()) return Futures.transformAsync(future, in -> in != null ? Futures.immediateFuture(in.getCustomerId())
: Futures.immediateFuture(null), MoreExecutors.directExecutor()); : Futures.immediateFuture(null), ctx.getDbCallbackExecutor());
} }
} }

3
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/util/EntitiesRelatedDeviceIdAsyncLoader.java

@ -17,7 +17,6 @@ package org.thingsboard.rule.engine.util;
import com.google.common.util.concurrent.Futures; import com.google.common.util.concurrent.Futures;
import com.google.common.util.concurrent.ListenableFuture; import com.google.common.util.concurrent.ListenableFuture;
import com.google.common.util.concurrent.MoreExecutors;
import org.apache.commons.collections.CollectionUtils; import org.apache.commons.collections.CollectionUtils;
import org.thingsboard.rule.engine.api.TbContext; import org.thingsboard.rule.engine.api.TbContext;
import org.thingsboard.rule.engine.data.DeviceRelationsQuery; import org.thingsboard.rule.engine.data.DeviceRelationsQuery;
@ -39,7 +38,7 @@ public class EntitiesRelatedDeviceIdAsyncLoader {
return Futures.transformAsync(devicesListFuture, return Futures.transformAsync(devicesListFuture,
deviceList -> CollectionUtils.isNotEmpty(deviceList) ? deviceList -> CollectionUtils.isNotEmpty(deviceList) ?
Futures.immediateFuture(deviceList.get(0).getId()) Futures.immediateFuture(deviceList.get(0).getId())
: Futures.immediateFuture(null), MoreExecutors.directExecutor()); : Futures.immediateFuture(null), ctx.getDbCallbackExecutor());
} }
private static DeviceSearchQuery buildQuery(EntityId originator, DeviceRelationsQuery deviceRelationsQuery) { private static DeviceSearchQuery buildQuery(EntityId originator, DeviceRelationsQuery deviceRelationsQuery) {

5
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/util/EntitiesRelatedEntityIdAsyncLoader.java

@ -17,7 +17,6 @@ package org.thingsboard.rule.engine.util;
import com.google.common.util.concurrent.Futures; import com.google.common.util.concurrent.Futures;
import com.google.common.util.concurrent.ListenableFuture; import com.google.common.util.concurrent.ListenableFuture;
import com.google.common.util.concurrent.MoreExecutors;
import org.apache.commons.collections.CollectionUtils; import org.apache.commons.collections.CollectionUtils;
import org.thingsboard.rule.engine.api.TbContext; import org.thingsboard.rule.engine.api.TbContext;
import org.thingsboard.rule.engine.data.RelationsQuery; import org.thingsboard.rule.engine.data.RelationsQuery;
@ -40,12 +39,12 @@ public class EntitiesRelatedEntityIdAsyncLoader {
return Futures.transformAsync(relationListFuture, return Futures.transformAsync(relationListFuture,
relationList -> CollectionUtils.isNotEmpty(relationList) ? relationList -> CollectionUtils.isNotEmpty(relationList) ?
Futures.immediateFuture(relationList.get(0).getTo()) Futures.immediateFuture(relationList.get(0).getTo())
: Futures.immediateFuture(null), MoreExecutors.directExecutor()); : Futures.immediateFuture(null), ctx.getDbCallbackExecutor());
} else if (relationsQuery.getDirection() == EntitySearchDirection.TO) { } else if (relationsQuery.getDirection() == EntitySearchDirection.TO) {
return Futures.transformAsync(relationListFuture, return Futures.transformAsync(relationListFuture,
relationList -> CollectionUtils.isNotEmpty(relationList) ? relationList -> CollectionUtils.isNotEmpty(relationList) ?
Futures.immediateFuture(relationList.get(0).getFrom()) Futures.immediateFuture(relationList.get(0).getFrom())
: Futures.immediateFuture(null), MoreExecutors.directExecutor()); : Futures.immediateFuture(null), ctx.getDbCallbackExecutor());
} }
return Futures.immediateFailedFuture(new IllegalStateException("Unknown direction")); return Futures.immediateFailedFuture(new IllegalStateException("Unknown direction"));
} }

38
rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/metadata/TbGetCustomerDetailsNodeTest.java

@ -16,6 +16,8 @@
package org.thingsboard.rule.engine.metadata; package org.thingsboard.rule.engine.metadata;
import com.google.common.util.concurrent.Futures; 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.BeforeEach;
import org.junit.jupiter.api.Test; import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.extension.ExtendWith; import org.junit.jupiter.api.extension.ExtendWith;
@ -23,6 +25,7 @@ import org.mockito.ArgumentCaptor;
import org.mockito.Mock; import org.mockito.Mock;
import org.mockito.junit.jupiter.MockitoExtension; import org.mockito.junit.jupiter.MockitoExtension;
import org.thingsboard.common.util.JacksonUtil; import org.thingsboard.common.util.JacksonUtil;
import org.thingsboard.common.util.ListeningExecutor;
import org.thingsboard.rule.engine.api.TbContext; import org.thingsboard.rule.engine.api.TbContext;
import org.thingsboard.rule.engine.api.TbNodeConfiguration; import org.thingsboard.rule.engine.api.TbNodeConfiguration;
import org.thingsboard.rule.engine.api.TbNodeException; import org.thingsboard.rule.engine.api.TbNodeException;
@ -55,6 +58,7 @@ import org.thingsboard.server.dao.user.UserService;
import java.util.Collections; import java.util.Collections;
import java.util.List; import java.util.List;
import java.util.UUID; import java.util.UUID;
import java.util.concurrent.Callable;
import static org.assertj.core.api.Assertions.assertThat; import static org.assertj.core.api.Assertions.assertThat;
import static org.junit.jupiter.api.Assertions.assertThrows; import static org.junit.jupiter.api.Assertions.assertThrows;
@ -71,6 +75,21 @@ public class TbGetCustomerDetailsNodeTest {
private static final DeviceId DUMMY_DEVICE_ORIGINATOR = new DeviceId(UUID.randomUUID()); private static final DeviceId DUMMY_DEVICE_ORIGINATOR = new DeviceId(UUID.randomUUID());
private static final TenantId TENANT_ID = new TenantId(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 @Mock
private TbContext ctxMock; private TbContext ctxMock;
@Mock @Mock
@ -209,6 +228,8 @@ public class TbGetCustomerDetailsNodeTest {
mockFindCustomer(); mockFindCustomer();
when(ctxMock.getDbCallbackExecutor()).thenReturn(DB_EXECUTOR);
// WHEN // WHEN
node.onMsg(ctxMock, msg); node.onMsg(ctxMock, msg);
@ -249,6 +270,8 @@ public class TbGetCustomerDetailsNodeTest {
mockFindCustomer(); mockFindCustomer();
when(ctxMock.getDbCallbackExecutor()).thenReturn(DB_EXECUTOR);
// WHEN // WHEN
node.onMsg(ctxMock, msg); node.onMsg(ctxMock, msg);
@ -281,6 +304,8 @@ public class TbGetCustomerDetailsNodeTest {
mockFindCustomer(); mockFindCustomer();
when(ctxMock.getDbCallbackExecutor()).thenReturn(DB_EXECUTOR);
// WHEN // WHEN
node.onMsg(ctxMock, msg); node.onMsg(ctxMock, msg);
@ -310,6 +335,8 @@ public class TbGetCustomerDetailsNodeTest {
when(ctxMock.getUserService()).thenReturn(userServiceMock); when(ctxMock.getUserService()).thenReturn(userServiceMock);
when(userServiceMock.findUserByIdAsync(eq(TENANT_ID), eq(user.getId()))).thenReturn(Futures.immediateFuture(user)); when(userServiceMock.findUserByIdAsync(eq(TENANT_ID), eq(user.getId()))).thenReturn(Futures.immediateFuture(user));
when(ctxMock.getDbCallbackExecutor()).thenReturn(DB_EXECUTOR);
mockFindCustomer(); mockFindCustomer();
// WHEN // WHEN
@ -334,13 +361,16 @@ public class TbGetCustomerDetailsNodeTest {
prepareMsgAndConfig(FetchTo.DATA, List.of(EntityDetails.ZIP, EntityDetails.ADDRESS, EntityDetails.ADDRESS2), edge.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(ctxMock.getEdgeService()).thenReturn(edgeServiceMock);
when(edgeServiceMock.findEdgeByIdAsync(eq(TENANT_ID), eq(edge.getId()))).thenReturn(Futures.immediateFuture(edge)); when(edgeServiceMock.findEdgeByIdAsync(eq(TENANT_ID), eq(edge.getId()))).thenReturn(Futures.immediateFuture(edge));
when(ctxMock.getTenantId()).thenReturn(TENANT_ID);
when(ctxMock.getCustomerService()).thenReturn(customerServiceMock); when(ctxMock.getCustomerService()).thenReturn(customerServiceMock);
when(customerServiceMock.findCustomerByIdAsync(eq(TENANT_ID), eq(customer.getId()))).thenReturn(Futures.immediateFuture(null)); when(customerServiceMock.findCustomerByIdAsync(eq(TENANT_ID), eq(customer.getId()))).thenReturn(Futures.immediateFuture(null));
when(ctxMock.getDbCallbackExecutor()).thenReturn(DB_EXECUTOR);
// WHEN // WHEN
node.onMsg(ctxMock, msg); node.onMsg(ctxMock, msg);
@ -368,6 +398,8 @@ public class TbGetCustomerDetailsNodeTest {
when(ctxMock.getEdgeService()).thenReturn(edgeServiceMock); when(ctxMock.getEdgeService()).thenReturn(edgeServiceMock);
when(edgeServiceMock.findEdgeByIdAsync(eq(TENANT_ID), eq(edge.getId()))).thenReturn(Futures.immediateFuture(null)); when(edgeServiceMock.findEdgeByIdAsync(eq(TENANT_ID), eq(edge.getId()))).thenReturn(Futures.immediateFuture(null));
when(ctxMock.getDbCallbackExecutor()).thenReturn(DB_EXECUTOR);
// WHEN // WHEN
node.onMsg(ctxMock, msg); node.onMsg(ctxMock, msg);
@ -394,6 +426,8 @@ public class TbGetCustomerDetailsNodeTest {
when(ctxMock.getDeviceService()).thenReturn(deviceServiceMock); when(ctxMock.getDeviceService()).thenReturn(deviceServiceMock);
when(deviceServiceMock.findDeviceByIdAsync(eq(TENANT_ID), eq(device.getId()))).thenReturn(Futures.immediateFuture(device)); when(deviceServiceMock.findDeviceByIdAsync(eq(TENANT_ID), eq(device.getId()))).thenReturn(Futures.immediateFuture(device));
when(ctxMock.getDbCallbackExecutor()).thenReturn(DB_EXECUTOR);
// WHEN // WHEN
node.onMsg(ctxMock, msg); node.onMsg(ctxMock, msg);
@ -421,6 +455,8 @@ public class TbGetCustomerDetailsNodeTest {
when(ctxMock.getDeviceService()).thenReturn(deviceServiceMock); when(ctxMock.getDeviceService()).thenReturn(deviceServiceMock);
when(deviceServiceMock.findDeviceByIdAsync(eq(TENANT_ID), eq(device.getId()))).thenReturn(Futures.immediateFuture(device)); when(deviceServiceMock.findDeviceByIdAsync(eq(TENANT_ID), eq(device.getId()))).thenReturn(Futures.immediateFuture(device));
when(ctxMock.getDbCallbackExecutor()).thenReturn(DB_EXECUTOR);
mockFindCustomer(); mockFindCustomer();
// WHEN // WHEN

31
rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/util/EntitiesCustomerIdAsyncLoaderTest.java

@ -16,11 +16,13 @@
package org.thingsboard.rule.engine.util; package org.thingsboard.rule.engine.util;
import com.google.common.util.concurrent.Futures; import com.google.common.util.concurrent.Futures;
import org.junit.jupiter.api.Assertions; import com.google.common.util.concurrent.ListenableFuture;
import org.jetbrains.annotations.NotNull;
import org.junit.jupiter.api.Test; import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.extension.ExtendWith; import org.junit.jupiter.api.extension.ExtendWith;
import org.mockito.Mock; import org.mockito.Mock;
import org.mockito.junit.jupiter.MockitoExtension; 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.TbContext;
import org.thingsboard.rule.engine.api.TbNodeException; import org.thingsboard.rule.engine.api.TbNodeException;
import org.thingsboard.server.common.data.Customer; import org.thingsboard.server.common.data.Customer;
@ -39,6 +41,7 @@ import org.thingsboard.server.dao.user.UserService;
import java.util.EnumSet; import java.util.EnumSet;
import java.util.UUID; import java.util.UUID;
import java.util.concurrent.Callable;
import java.util.concurrent.ExecutionException; import java.util.concurrent.ExecutionException;
import static org.junit.jupiter.api.Assertions.assertEquals; import static org.junit.jupiter.api.Assertions.assertEquals;
@ -57,6 +60,21 @@ public class EntitiesCustomerIdAsyncLoaderTest {
EntityType.ASSET, EntityType.ASSET,
EntityType.DEVICE 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 @Mock
private TbContext ctxMock; private TbContext ctxMock;
@Mock @Mock
@ -75,7 +93,7 @@ public class EntitiesCustomerIdAsyncLoaderTest {
var actualCustomerId = EntitiesCustomerIdAsyncLoader.findEntityIdAsync(ctxMock, customer.getId()).get(); var actualCustomerId = EntitiesCustomerIdAsyncLoader.findEntityIdAsync(ctxMock, customer.getId()).get();
// THEN // THEN
Assertions.assertEquals(customer.getId(), actualCustomerId); assertEquals(customer.getId(), actualCustomerId);
} }
@Test @Test
@ -87,12 +105,13 @@ public class EntitiesCustomerIdAsyncLoaderTest {
when(ctxMock.getUserService()).thenReturn(userServiceMock); when(ctxMock.getUserService()).thenReturn(userServiceMock);
doReturn(Futures.immediateFuture(user)).when(userServiceMock).findUserByIdAsync(any(), any()); doReturn(Futures.immediateFuture(user)).when(userServiceMock).findUserByIdAsync(any(), any());
when(ctxMock.getDbCallbackExecutor()).thenReturn(DB_EXECUTOR);
// WHEN // WHEN
var actualCustomerId = EntitiesCustomerIdAsyncLoader.findEntityIdAsync(ctxMock, user.getId()).get(); var actualCustomerId = EntitiesCustomerIdAsyncLoader.findEntityIdAsync(ctxMock, user.getId()).get();
// THEN // THEN
Assertions.assertEquals(expectedCustomerId, actualCustomerId); assertEquals(expectedCustomerId, actualCustomerId);
} }
@Test @Test
@ -104,12 +123,13 @@ public class EntitiesCustomerIdAsyncLoaderTest {
when(ctxMock.getAssetService()).thenReturn(assetServiceMock); when(ctxMock.getAssetService()).thenReturn(assetServiceMock);
doReturn(Futures.immediateFuture(asset)).when(assetServiceMock).findAssetByIdAsync(any(), any()); doReturn(Futures.immediateFuture(asset)).when(assetServiceMock).findAssetByIdAsync(any(), any());
when(ctxMock.getDbCallbackExecutor()).thenReturn(DB_EXECUTOR);
// WHEN // WHEN
var actualCustomerId = EntitiesCustomerIdAsyncLoader.findEntityIdAsync(ctxMock, asset.getId()).get(); var actualCustomerId = EntitiesCustomerIdAsyncLoader.findEntityIdAsync(ctxMock, asset.getId()).get();
// THEN // THEN
Assertions.assertEquals(expectedCustomerId, actualCustomerId); assertEquals(expectedCustomerId, actualCustomerId);
} }
@Test @Test
@ -121,12 +141,13 @@ public class EntitiesCustomerIdAsyncLoaderTest {
when(ctxMock.getDeviceService()).thenReturn(deviceServiceMock); when(ctxMock.getDeviceService()).thenReturn(deviceServiceMock);
doReturn(Futures.immediateFuture(device)).when(deviceServiceMock).findDeviceByIdAsync(any(), any()); doReturn(Futures.immediateFuture(device)).when(deviceServiceMock).findDeviceByIdAsync(any(), any());
when(ctxMock.getDbCallbackExecutor()).thenReturn(DB_EXECUTOR);
// WHEN // WHEN
var actualCustomerId = EntitiesCustomerIdAsyncLoader.findEntityIdAsync(ctxMock, device.getId()).get(); var actualCustomerId = EntitiesCustomerIdAsyncLoader.findEntityIdAsync(ctxMock, device.getId()).get();
// THEN // THEN
Assertions.assertEquals(expectedCustomerId, actualCustomerId); assertEquals(expectedCustomerId, actualCustomerId);
} }
@Test @Test

21
rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/util/EntitiesRelatedDeviceIdAsyncLoaderTest.java

@ -16,10 +16,13 @@
package org.thingsboard.rule.engine.util; package org.thingsboard.rule.engine.util;
import com.google.common.util.concurrent.Futures; 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.Test;
import org.junit.jupiter.api.extension.ExtendWith; import org.junit.jupiter.api.extension.ExtendWith;
import org.mockito.Mock; import org.mockito.Mock;
import org.mockito.junit.jupiter.MockitoExtension; 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.TbContext;
import org.thingsboard.rule.engine.data.DeviceRelationsQuery; import org.thingsboard.rule.engine.data.DeviceRelationsQuery;
import org.thingsboard.server.common.data.Device; import org.thingsboard.server.common.data.Device;
@ -34,6 +37,7 @@ import org.thingsboard.server.dao.device.DeviceService;
import java.util.List; import java.util.List;
import java.util.UUID; import java.util.UUID;
import java.util.concurrent.Callable;
import static org.junit.jupiter.api.Assertions.assertEquals; import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertNotNull; import static org.junit.jupiter.api.Assertions.assertNotNull;
@ -47,6 +51,21 @@ public class EntitiesRelatedDeviceIdAsyncLoaderTest {
private static final EntityId DUMMY_ORIGINATOR = new DeviceId(UUID.randomUUID()); private static final EntityId DUMMY_ORIGINATOR = new DeviceId(UUID.randomUUID());
private static final TenantId TENANT_ID = new TenantId(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 @Mock
private TbContext ctxMock; private TbContext ctxMock;
@Mock @Mock
@ -76,6 +95,7 @@ public class EntitiesRelatedDeviceIdAsyncLoaderTest {
when(ctxMock.getDeviceService()).thenReturn(deviceServiceMock); when(ctxMock.getDeviceService()).thenReturn(deviceServiceMock);
when(deviceServiceMock.findDevicesByQuery(eq(TENANT_ID), eq(expectedDeviceSearchQuery))) when(deviceServiceMock.findDevicesByQuery(eq(TENANT_ID), eq(expectedDeviceSearchQuery)))
.thenReturn(Futures.immediateFuture(null)); .thenReturn(Futures.immediateFuture(null));
when(ctxMock.getDbCallbackExecutor()).thenReturn(DB_EXECUTOR);
// WHEN // WHEN
EntitiesRelatedDeviceIdAsyncLoader.findDeviceAsync(ctxMock, DUMMY_ORIGINATOR, deviceRelationsQuery); EntitiesRelatedDeviceIdAsyncLoader.findDeviceAsync(ctxMock, DUMMY_ORIGINATOR, deviceRelationsQuery);
@ -117,6 +137,7 @@ public class EntitiesRelatedDeviceIdAsyncLoaderTest {
when(ctxMock.getDeviceService()).thenReturn(deviceServiceMock); when(ctxMock.getDeviceService()).thenReturn(deviceServiceMock);
when(deviceServiceMock.findDevicesByQuery(eq(TENANT_ID), eq(expectedDeviceSearchQuery))) when(deviceServiceMock.findDevicesByQuery(eq(TENANT_ID), eq(expectedDeviceSearchQuery)))
.thenReturn(Futures.immediateFuture(devicesList)); .thenReturn(Futures.immediateFuture(devicesList));
when(ctxMock.getDbCallbackExecutor()).thenReturn(DB_EXECUTOR);
// WHEN // WHEN
var entityIdFuture = EntitiesRelatedDeviceIdAsyncLoader.findDeviceAsync(ctxMock, DUMMY_ORIGINATOR, deviceRelationsQuery); var entityIdFuture = EntitiesRelatedDeviceIdAsyncLoader.findDeviceAsync(ctxMock, DUMMY_ORIGINATOR, deviceRelationsQuery);

21
rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/util/EntitiesRelatedEntitiesIdAsyncLoaderTest.java

@ -16,10 +16,13 @@
package org.thingsboard.rule.engine.util; package org.thingsboard.rule.engine.util;
import com.google.common.util.concurrent.Futures; 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.Test;
import org.junit.jupiter.api.extension.ExtendWith; import org.junit.jupiter.api.extension.ExtendWith;
import org.mockito.Mock; import org.mockito.Mock;
import org.mockito.junit.jupiter.MockitoExtension; 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.TbContext;
import org.thingsboard.rule.engine.data.RelationsQuery; import org.thingsboard.rule.engine.data.RelationsQuery;
import org.thingsboard.server.common.data.Device; import org.thingsboard.server.common.data.Device;
@ -39,6 +42,7 @@ import org.thingsboard.server.dao.relation.RelationService;
import java.util.Collections; import java.util.Collections;
import java.util.List; import java.util.List;
import java.util.UUID; import java.util.UUID;
import java.util.concurrent.Callable;
import static org.junit.jupiter.api.Assertions.assertEquals; import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertNotNull; import static org.junit.jupiter.api.Assertions.assertNotNull;
@ -52,6 +56,21 @@ public class EntitiesRelatedEntitiesIdAsyncLoaderTest {
private static final EntityId DUMMY_ORIGINATOR = new DeviceId(UUID.randomUUID()); private static final EntityId DUMMY_ORIGINATOR = new DeviceId(UUID.randomUUID());
private static final TenantId TENANT_ID = new TenantId(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 @Mock
private TbContext ctxMock; private TbContext ctxMock;
@Mock @Mock
@ -80,6 +99,7 @@ public class EntitiesRelatedEntitiesIdAsyncLoaderTest {
when(ctxMock.getRelationService()).thenReturn(relationServiceMock); when(ctxMock.getRelationService()).thenReturn(relationServiceMock);
when(relationServiceMock.findByQuery(eq(TENANT_ID), eq(expectedEntityRelationsQuery))) when(relationServiceMock.findByQuery(eq(TENANT_ID), eq(expectedEntityRelationsQuery)))
.thenReturn(Futures.immediateFuture(null)); .thenReturn(Futures.immediateFuture(null));
when(ctxMock.getDbCallbackExecutor()).thenReturn(DB_EXECUTOR);
// WHEN // WHEN
EntitiesRelatedEntityIdAsyncLoader.findEntityAsync(ctxMock, DUMMY_ORIGINATOR, relationsQuery); EntitiesRelatedEntityIdAsyncLoader.findEntityAsync(ctxMock, DUMMY_ORIGINATOR, relationsQuery);
@ -138,6 +158,7 @@ public class EntitiesRelatedEntitiesIdAsyncLoaderTest {
when(ctxMock.getRelationService()).thenReturn(relationServiceMock); when(ctxMock.getRelationService()).thenReturn(relationServiceMock);
when(relationServiceMock.findByQuery(eq(TENANT_ID), eq(expectedEntityRelationsQuery))) when(relationServiceMock.findByQuery(eq(TENANT_ID), eq(expectedEntityRelationsQuery)))
.thenReturn(Futures.immediateFuture(expectedEntityRelationsList)); .thenReturn(Futures.immediateFuture(expectedEntityRelationsList));
when(ctxMock.getDbCallbackExecutor()).thenReturn(DB_EXECUTOR);
// WHEN // WHEN
var deviceIdFuture = EntitiesRelatedEntityIdAsyncLoader.findEntityAsync(ctxMock, DUMMY_ORIGINATOR, relationsQuery); var deviceIdFuture = EntitiesRelatedEntityIdAsyncLoader.findEntityAsync(ctxMock, DUMMY_ORIGINATOR, relationsQuery);

Loading…
Cancel
Save