From 913c6b60e234f1d29fd42d95887e05a28ebebd49 Mon Sep 17 00:00:00 2001 From: Dmytro Skarzhynets Date: Mon, 17 Apr 2023 14:41:56 +0300 Subject: [PATCH] Review fixes: replace direct executor with db executor --- .../metadata/TbAbstractGetAttributesNode.java | 6 +-- .../metadata/TbAbstractGetEntityAttrNode.java | 7 ++-- .../metadata/TbGetCustomerDetailsNode.java | 11 +++--- .../util/EntitiesCustomerIdAsyncLoader.java | 11 +++--- .../EntitiesRelatedDeviceIdAsyncLoader.java | 3 +- .../EntitiesRelatedEntityIdAsyncLoader.java | 5 +-- .../TbGetCustomerDetailsNodeTest.java | 38 ++++++++++++++++++- .../EntitiesCustomerIdAsyncLoaderTest.java | 31 ++++++++++++--- ...ntitiesRelatedDeviceIdAsyncLoaderTest.java | 21 ++++++++++ ...itiesRelatedEntitiesIdAsyncLoaderTest.java | 21 ++++++++++ 10 files changed, 125 insertions(+), 29 deletions(-) diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbAbstractGetAttributesNode.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbAbstractGetAttributesNode.java index b88ae1a0a1..43ee96c917 100644 --- a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbAbstractGetAttributesNode.java +++ b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbAbstractGetAttributesNode.java @@ -101,7 +101,7 @@ public abstract class TbAbstractGetAttributesNode ctx.tellFailure(msg, t), ctx.getDbCallbackExecutor()); + }, t -> ctx.tellFailure(msg, t), MoreExecutors.directExecutor()); } private ListenableFuture>> getAttrAsync( @@ -121,7 +121,7 @@ public abstract class TbAbstractGetAttributesNode(scope, nonExistentKeys)); } return new TbPair<>(scope, attributeKvEntryList); - }, MoreExecutors.directExecutor()); + }, ctx.getDbCallbackExecutor()); } private ListenableFuture>> getLatestTelemetry(TbContext ctx, EntityId entityId, List keys, Set>> failuresPairSet) { @@ -147,7 +147,7 @@ public abstract class TbAbstractGetAttributesNode(LATEST_TS, nonExistentKeys)); } return new TbPair<>(LATEST_TS, listTsKvEntry); - }, MoreExecutors.directExecutor()); + }, ctx.getDbCallbackExecutor()); } private TsKvEntry getValueWithTs(TsKvEntry tsKvEntry) { diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbAbstractGetEntityAttrNode.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbAbstractGetEntityAttrNode.java index c11e6b8c2a..2b04d6e6e7 100644 --- a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbAbstractGetEntityAttrNode.java +++ b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbAbstractGetEntityAttrNode.java @@ -65,7 +65,8 @@ public abstract class TbAbstractGetEntityAttrNode extends Tb var sourceKeys = List.copyOf(mappingsMap.keySet()); withCallback(config.isTelemetry() ? getLatestTelemetryAsync(ctx, entityId, sourceKeys) : getAttributesAsync(ctx, entityId, sourceKeys), data -> putDataAndTell(ctx, msg, data, mappingsMap, msgDataAsJsonNode), - t -> ctx.tellFailure(msg, t), ctx.getDbCallbackExecutor()); + t -> ctx.tellFailure(msg, t), + MoreExecutors.directExecutor()); } private ListenableFuture> getAttributesAsync(TbContext ctx, EntityId entityId, List attrKeys) { @@ -74,7 +75,7 @@ public abstract class TbAbstractGetEntityAttrNode extends Tb l.stream() .map(i -> (KvEntry) i) .collect(Collectors.toList()), - MoreExecutors.directExecutor()); + ctx.getDbCallbackExecutor()); } private ListenableFuture> getLatestTelemetryAsync(TbContext ctx, EntityId entityId, List timeseriesKeys) { @@ -83,7 +84,7 @@ public abstract class TbAbstractGetEntityAttrNode extends Tb l.stream() .map(i -> (KvEntry) i) .collect(Collectors.toList()), - MoreExecutors.directExecutor()); + ctx.getDbCallbackExecutor()); } private void putDataAndTell(TbContext ctx, TbMsg msg, List data, Map map, ObjectNode msgData) { diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetCustomerDetailsNode.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetCustomerDetailsNode.java index 2c7fe2302b..65f11facf5 100644 --- a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetCustomerDetailsNode.java +++ b/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.ListenableFuture; -import com.google.common.util.concurrent.MoreExecutors; import lombok.extern.slf4j.Slf4j; import org.thingsboard.rule.engine.api.RuleNode; import org.thingsboard.rule.engine.api.TbContext; @@ -70,19 +69,19 @@ public class TbGetCustomerDetailsNode extends TbAbstractGetEntityDetailsNode getCustomerFuture(ctx, device, msg.getOriginator()), MoreExecutors.directExecutor()); + device -> getCustomerFuture(ctx, device, msg.getOriginator()), ctx.getDbCallbackExecutor()); case ASSET: 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: 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: 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: 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: return Futures.immediateFailedFuture(new NoSuchElementException("Entity with entityType '" + msg.getOriginator().getEntityType() + "' is not supported.")); } diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/util/EntitiesCustomerIdAsyncLoader.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/util/EntitiesCustomerIdAsyncLoader.java index d517afaf35..c5b6c0771e 100644 --- a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/util/EntitiesCustomerIdAsyncLoader.java +++ b/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.ListenableFuture; -import com.google.common.util.concurrent.MoreExecutors; import org.thingsboard.rule.engine.api.TbContext; import org.thingsboard.rule.engine.api.TbNodeException; import org.thingsboard.server.common.data.HasCustomerId; @@ -34,19 +33,19 @@ public class EntitiesCustomerIdAsyncLoader { case CUSTOMER: return Futures.immediateFuture((CustomerId) originator); case USER: - return toCustomerIdAsync(ctx.getUserService().findUserByIdAsync(ctx.getTenantId(), (UserId) originator)); + return toCustomerIdAsync(ctx, ctx.getUserService().findUserByIdAsync(ctx.getTenantId(), (UserId) originator)); case ASSET: - return toCustomerIdAsync(ctx.getAssetService().findAssetByIdAsync(ctx.getTenantId(), (AssetId) originator)); + return toCustomerIdAsync(ctx, ctx.getAssetService().findAssetByIdAsync(ctx.getTenantId(), (AssetId) originator)); case DEVICE: - return toCustomerIdAsync(ctx.getDeviceService().findDeviceByIdAsync(ctx.getTenantId(), (DeviceId) originator)); + return toCustomerIdAsync(ctx, ctx.getDeviceService().findDeviceByIdAsync(ctx.getTenantId(), (DeviceId) originator)); default: return Futures.immediateFailedFuture(new TbNodeException("Unexpected originator EntityType: " + originator.getEntityType())); } } - private static ListenableFuture toCustomerIdAsync(ListenableFuture future) { + private static ListenableFuture toCustomerIdAsync(TbContext ctx, ListenableFuture future) { return Futures.transformAsync(future, in -> in != null ? Futures.immediateFuture(in.getCustomerId()) - : Futures.immediateFuture(null), MoreExecutors.directExecutor()); + : Futures.immediateFuture(null), ctx.getDbCallbackExecutor()); } } diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/util/EntitiesRelatedDeviceIdAsyncLoader.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/util/EntitiesRelatedDeviceIdAsyncLoader.java index 1a590cc32f..b937f01b24 100644 --- a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/util/EntitiesRelatedDeviceIdAsyncLoader.java +++ b/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.ListenableFuture; -import com.google.common.util.concurrent.MoreExecutors; import org.apache.commons.collections.CollectionUtils; import org.thingsboard.rule.engine.api.TbContext; import org.thingsboard.rule.engine.data.DeviceRelationsQuery; @@ -39,7 +38,7 @@ public class EntitiesRelatedDeviceIdAsyncLoader { return Futures.transformAsync(devicesListFuture, deviceList -> CollectionUtils.isNotEmpty(deviceList) ? Futures.immediateFuture(deviceList.get(0).getId()) - : Futures.immediateFuture(null), MoreExecutors.directExecutor()); + : Futures.immediateFuture(null), ctx.getDbCallbackExecutor()); } private static DeviceSearchQuery buildQuery(EntityId originator, DeviceRelationsQuery deviceRelationsQuery) { diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/util/EntitiesRelatedEntityIdAsyncLoader.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/util/EntitiesRelatedEntityIdAsyncLoader.java index dbad7c9c9d..d0caac8876 100644 --- a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/util/EntitiesRelatedEntityIdAsyncLoader.java +++ b/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.ListenableFuture; -import com.google.common.util.concurrent.MoreExecutors; import org.apache.commons.collections.CollectionUtils; import org.thingsboard.rule.engine.api.TbContext; import org.thingsboard.rule.engine.data.RelationsQuery; @@ -40,12 +39,12 @@ public class EntitiesRelatedEntityIdAsyncLoader { return Futures.transformAsync(relationListFuture, relationList -> CollectionUtils.isNotEmpty(relationList) ? Futures.immediateFuture(relationList.get(0).getTo()) - : Futures.immediateFuture(null), MoreExecutors.directExecutor()); + : Futures.immediateFuture(null), ctx.getDbCallbackExecutor()); } else if (relationsQuery.getDirection() == EntitySearchDirection.TO) { return Futures.transformAsync(relationListFuture, relationList -> CollectionUtils.isNotEmpty(relationList) ? Futures.immediateFuture(relationList.get(0).getFrom()) - : Futures.immediateFuture(null), MoreExecutors.directExecutor()); + : Futures.immediateFuture(null), ctx.getDbCallbackExecutor()); } return Futures.immediateFailedFuture(new IllegalStateException("Unknown direction")); } diff --git a/rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/metadata/TbGetCustomerDetailsNodeTest.java b/rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/metadata/TbGetCustomerDetailsNodeTest.java index 0c78d8e346..1cb1449078 100644 --- a/rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/metadata/TbGetCustomerDetailsNodeTest.java +++ b/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; 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; @@ -23,6 +25,7 @@ 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; @@ -55,6 +58,7 @@ import org.thingsboard.server.dao.user.UserService; import java.util.Collections; import java.util.List; 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; @@ -71,6 +75,21 @@ 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 ListenableFuture executeAsync(Callable 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 @@ -209,6 +228,8 @@ public class TbGetCustomerDetailsNodeTest { mockFindCustomer(); + when(ctxMock.getDbCallbackExecutor()).thenReturn(DB_EXECUTOR); + // WHEN node.onMsg(ctxMock, msg); @@ -249,6 +270,8 @@ public class TbGetCustomerDetailsNodeTest { mockFindCustomer(); + when(ctxMock.getDbCallbackExecutor()).thenReturn(DB_EXECUTOR); + // WHEN node.onMsg(ctxMock, msg); @@ -281,6 +304,8 @@ public class TbGetCustomerDetailsNodeTest { mockFindCustomer(); + when(ctxMock.getDbCallbackExecutor()).thenReturn(DB_EXECUTOR); + // WHEN node.onMsg(ctxMock, msg); @@ -310,6 +335,8 @@ public class TbGetCustomerDetailsNodeTest { 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 @@ -334,13 +361,16 @@ public class TbGetCustomerDetailsNodeTest { 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.getTenantId()).thenReturn(TENANT_ID); 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); @@ -368,6 +398,8 @@ public class TbGetCustomerDetailsNodeTest { 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); @@ -394,6 +426,8 @@ public class TbGetCustomerDetailsNodeTest { 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); @@ -421,6 +455,8 @@ public class TbGetCustomerDetailsNodeTest { 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 diff --git a/rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/util/EntitiesCustomerIdAsyncLoaderTest.java b/rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/util/EntitiesCustomerIdAsyncLoaderTest.java index 97201daeda..2edc41f936 100644 --- a/rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/util/EntitiesCustomerIdAsyncLoaderTest.java +++ b/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; 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.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; @@ -39,6 +41,7 @@ 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; @@ -57,6 +60,21 @@ public class EntitiesCustomerIdAsyncLoaderTest { EntityType.ASSET, EntityType.DEVICE ); + private static final ListeningExecutor DB_EXECUTOR = new ListeningExecutor() { + @Override + public ListenableFuture executeAsync(Callable 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 @@ -75,7 +93,7 @@ public class EntitiesCustomerIdAsyncLoaderTest { var actualCustomerId = EntitiesCustomerIdAsyncLoader.findEntityIdAsync(ctxMock, customer.getId()).get(); // THEN - Assertions.assertEquals(customer.getId(), actualCustomerId); + assertEquals(customer.getId(), actualCustomerId); } @Test @@ -87,12 +105,13 @@ public class EntitiesCustomerIdAsyncLoaderTest { 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 - Assertions.assertEquals(expectedCustomerId, actualCustomerId); + assertEquals(expectedCustomerId, actualCustomerId); } @Test @@ -104,12 +123,13 @@ public class EntitiesCustomerIdAsyncLoaderTest { 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 - Assertions.assertEquals(expectedCustomerId, actualCustomerId); + assertEquals(expectedCustomerId, actualCustomerId); } @Test @@ -121,12 +141,13 @@ public class EntitiesCustomerIdAsyncLoaderTest { 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 - Assertions.assertEquals(expectedCustomerId, actualCustomerId); + assertEquals(expectedCustomerId, actualCustomerId); } @Test diff --git a/rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/util/EntitiesRelatedDeviceIdAsyncLoaderTest.java b/rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/util/EntitiesRelatedDeviceIdAsyncLoaderTest.java index fa21f3746b..1c0faec7bc 100644 --- a/rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/util/EntitiesRelatedDeviceIdAsyncLoaderTest.java +++ b/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; 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; @@ -34,6 +37,7 @@ 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; @@ -47,6 +51,21 @@ 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 ListenableFuture executeAsync(Callable 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 @@ -76,6 +95,7 @@ public class EntitiesRelatedDeviceIdAsyncLoaderTest { 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); @@ -117,6 +137,7 @@ public class EntitiesRelatedDeviceIdAsyncLoaderTest { 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); diff --git a/rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/util/EntitiesRelatedEntitiesIdAsyncLoaderTest.java b/rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/util/EntitiesRelatedEntitiesIdAsyncLoaderTest.java index aae13a52a9..fc93206c18 100644 --- a/rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/util/EntitiesRelatedEntitiesIdAsyncLoaderTest.java +++ b/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; 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; @@ -39,6 +42,7 @@ 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; @@ -52,6 +56,21 @@ 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 ListenableFuture executeAsync(Callable 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 @@ -80,6 +99,7 @@ public class EntitiesRelatedEntitiesIdAsyncLoaderTest { 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); @@ -138,6 +158,7 @@ public class EntitiesRelatedEntitiesIdAsyncLoaderTest { 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);