From 9a34396d4521a553c56b590714ed49fc2e62fa12 Mon Sep 17 00:00:00 2001 From: Dmytro Skarzhynets Date: Wed, 29 Oct 2025 16:40:32 +0200 Subject: [PATCH] Support all customer-owned entities in 'customer attributes' and 'change originator' rule nodes --- .../server/dao/entity/EntityDaoService.java | 3 + .../server/dao/entity/EntityService.java | 4 + .../dao/entityview/EntityViewService.java | 8 +- .../server/common/data/HasCustomerId.java | 1 + .../server/dao/ai/AiModelServiceImpl.java | 7 + .../server/dao/alarm/BaseAlarmService.java | 10 +- .../dao/asset/AssetProfileServiceImpl.java | 8 + .../server/dao/asset/BaseAssetService.java | 11 +- .../dao/cf/BaseCalculatedFieldService.java | 8 + .../dao/customer/CustomerServiceImpl.java | 8 + .../dao/dashboard/DashboardServiceImpl.java | 9 + .../dao/device/DeviceProfileServiceImpl.java | 11 +- .../server/dao/device/DeviceServiceImpl.java | 26 +-- .../server/dao/domain/DomainServiceImpl.java | 16 +- .../server/dao/edge/EdgeServiceImpl.java | 10 +- .../server/dao/entity/BaseEntityService.java | 20 ++- .../dao/entityview/EntityViewServiceImpl.java | 30 ++-- .../server/dao/job/DefaultJobService.java | 8 + .../mobile/MobileAppBundleServiceImpl.java | 15 +- .../dao/mobile/MobileAppServiceImpl.java | 14 +- .../DefaultNotificationRequestService.java | 13 +- .../DefaultNotificationRuleService.java | 14 +- .../DefaultNotificationService.java | 9 + .../DefaultNotificationTargetService.java | 14 +- .../DefaultNotificationTemplateService.java | 9 + .../dao/oauth2/OAuth2ClientServiceImpl.java | 9 +- .../server/dao/ota/BaseOtaPackageService.java | 8 + .../server/dao/queue/BaseQueueService.java | 36 ++-- .../dao/queue/BaseQueueStatsService.java | 8 + .../server/dao/resource/BaseImageService.java | 34 ++-- .../dao/resource/BaseResourceService.java | 13 +- .../server/dao/rpc/BaseRpcService.java | 35 ++-- .../server/dao/rule/BaseRuleChainService.java | 38 +++-- .../settings/AdminSettingsServiceImpl.java | 9 + .../dao/tenant/TenantProfileServiceImpl.java | 30 ++-- .../server/dao/tenant/TenantServiceImpl.java | 10 +- .../usagerecord/ApiUsageStateServiceImpl.java | 11 +- .../server/dao/user/UserServiceImpl.java | 10 +- .../dao/widget/WidgetTypeServiceImpl.java | 12 +- .../dao/widget/WidgetsBundleServiceImpl.java | 48 +++--- .../metadata/TbAbstractGetEntityDataNode.java | 12 +- .../metadata/TbAbstractNodeWithFetchTo.java | 9 +- .../metadata/TbGetCustomerAttributeNode.java | 26 ++- .../engine/metadata/TbGetDeviceAttrNode.java | 2 +- .../metadata/TbGetRelatedAttributeNode.java | 2 +- .../transform/TbChangeOriginatorNode.java | 14 +- .../util/EntitiesCustomerIdAsyncLoader.java | 50 ------ .../rule/engine/TestDbCallbackExecutor.java | 2 +- .../TbGetCustomerAttributeNodeTest.java | 111 ++++++------- .../transform/TbChangeOriginatorNodeTest.java | 90 +++++++--- .../EntitiesCustomerIdAsyncLoaderTest.java | 154 ------------------ 51 files changed, 585 insertions(+), 484 deletions(-) delete mode 100644 rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/util/EntitiesCustomerIdAsyncLoader.java delete mode 100644 rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/util/EntitiesCustomerIdAsyncLoaderTest.java diff --git a/common/dao-api/src/main/java/org/thingsboard/server/dao/entity/EntityDaoService.java b/common/dao-api/src/main/java/org/thingsboard/server/dao/entity/EntityDaoService.java index e9350842f8..dca2e81f29 100644 --- a/common/dao-api/src/main/java/org/thingsboard/server/dao/entity/EntityDaoService.java +++ b/common/dao-api/src/main/java/org/thingsboard/server/dao/entity/EntityDaoService.java @@ -15,6 +15,7 @@ */ package org.thingsboard.server.dao.entity; +import com.google.common.util.concurrent.FluentFuture; import org.thingsboard.server.common.data.EntityType; import org.thingsboard.server.common.data.id.EntityId; import org.thingsboard.server.common.data.id.HasId; @@ -26,6 +27,8 @@ public interface EntityDaoService { Optional> findEntity(TenantId tenantId, EntityId entityId); + FluentFuture>> findEntityAsync(TenantId tenantId, EntityId entityId); + default long countByTenantId(TenantId tenantId) { throw new IllegalArgumentException("Not implemented for " + getEntityType()); } diff --git a/common/dao-api/src/main/java/org/thingsboard/server/dao/entity/EntityService.java b/common/dao-api/src/main/java/org/thingsboard/server/dao/entity/EntityService.java index 9adf703e0b..0e41e0b031 100644 --- a/common/dao-api/src/main/java/org/thingsboard/server/dao/entity/EntityService.java +++ b/common/dao-api/src/main/java/org/thingsboard/server/dao/entity/EntityService.java @@ -15,6 +15,7 @@ */ package org.thingsboard.server.dao.entity; +import com.google.common.util.concurrent.FluentFuture; import org.thingsboard.server.common.data.EntityInfo; import org.thingsboard.server.common.data.id.CustomerId; import org.thingsboard.server.common.data.id.EntityId; @@ -38,6 +39,8 @@ public interface EntityService { Optional fetchEntityCustomerId(TenantId tenantId, EntityId entityId); + FluentFuture> fetchEntityCustomerIdAsync(TenantId tenantId, EntityId entityId); + Optional> fetchEntity(TenantId tenantId, EntityId entityId); Map fetchEntityInfos(TenantId tenantId, CustomerId customerId, Set entityIds); @@ -47,4 +50,5 @@ public interface EntityService { long countEntitiesByQuery(TenantId tenantId, CustomerId customerId, EntityCountQuery query); PageData findEntityDataByQuery(TenantId tenantId, CustomerId customerId, EntityDataQuery query); + } diff --git a/common/dao-api/src/main/java/org/thingsboard/server/dao/entityview/EntityViewService.java b/common/dao-api/src/main/java/org/thingsboard/server/dao/entityview/EntityViewService.java index 6e557106f8..5ec2dfc20e 100644 --- a/common/dao-api/src/main/java/org/thingsboard/server/dao/entityview/EntityViewService.java +++ b/common/dao-api/src/main/java/org/thingsboard/server/dao/entityview/EntityViewService.java @@ -31,9 +31,6 @@ import org.thingsboard.server.dao.entity.EntityDaoService; import java.util.List; -/** - * Created by Victor Basanets on 8/27/2017. - */ public interface EntityViewService extends EntityDaoService { EntityView saveEntityView(EntityView entityView); @@ -52,6 +49,8 @@ public interface EntityViewService extends EntityDaoService { EntityView findEntityViewById(TenantId tenantId, EntityViewId entityViewId, boolean putInCache); + ListenableFuture findEntityViewByIdAsync(TenantId tenantId, EntityViewId entityViewId); + EntityView findEntityViewByTenantIdAndName(TenantId tenantId, String name); ListenableFuture findEntityViewByTenantIdAndNameAsync(TenantId tenantId, String name); @@ -74,8 +73,6 @@ public interface EntityViewService extends EntityDaoService { ListenableFuture> findEntityViewsByQuery(TenantId tenantId, EntityViewSearchQuery query); - ListenableFuture findEntityViewByIdAsync(TenantId tenantId, EntityViewId entityViewId); - ListenableFuture> findEntityViewsByTenantIdAndEntityIdAsync(TenantId tenantId, EntityId entityId); List findEntityViewsByTenantIdAndEntityId(TenantId tenantId, EntityId entityId); @@ -95,4 +92,5 @@ public interface EntityViewService extends EntityDaoService { PageData findEntityViewsByTenantIdAndEdgeId(TenantId tenantId, EdgeId edgeId, PageLink pageLink); PageData findEntityViewsByTenantIdAndEdgeIdAndType(TenantId tenantId, EdgeId edgeId, String type, PageLink pageLink); + } diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/HasCustomerId.java b/common/data/src/main/java/org/thingsboard/server/common/data/HasCustomerId.java index 6a5501840d..d7900e96ed 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/HasCustomerId.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/HasCustomerId.java @@ -20,4 +20,5 @@ import org.thingsboard.server.common.data.id.CustomerId; public interface HasCustomerId { CustomerId getCustomerId(); + } diff --git a/dao/src/main/java/org/thingsboard/server/dao/ai/AiModelServiceImpl.java b/dao/src/main/java/org/thingsboard/server/dao/ai/AiModelServiceImpl.java index b091a29247..102bf76ee6 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/ai/AiModelServiceImpl.java +++ b/dao/src/main/java/org/thingsboard/server/dao/ai/AiModelServiceImpl.java @@ -36,6 +36,7 @@ import org.thingsboard.server.dao.sql.JpaExecutorService; import java.util.Optional; import java.util.Set; +import static com.google.common.util.concurrent.MoreExecutors.directExecutor; import static org.thingsboard.server.dao.service.Validator.validatePageLink; @Service @@ -115,6 +116,12 @@ class AiModelServiceImpl extends CachedVersionedEntityService model); // necessary to cast to HasId } + @Override + public FluentFuture>> findEntityAsync(TenantId tenantId, EntityId entityId) { + return findAiModelByTenantIdAndIdAsync(tenantId, new AiModelId(entityId.getId())) + .transform(modelOpt -> modelOpt.map(model -> model), directExecutor()); // necessary to cast to HasId + } + @Override public long countByTenantId(TenantId tenantId) { return aiModelDao.countByTenantId(tenantId); diff --git a/dao/src/main/java/org/thingsboard/server/dao/alarm/BaseAlarmService.java b/dao/src/main/java/org/thingsboard/server/dao/alarm/BaseAlarmService.java index 5c3e4f3c79..9696a4cf31 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/alarm/BaseAlarmService.java +++ b/dao/src/main/java/org/thingsboard/server/dao/alarm/BaseAlarmService.java @@ -15,7 +15,6 @@ */ package org.thingsboard.server.dao.alarm; - import com.fasterxml.jackson.databind.JsonNode; import com.google.common.util.concurrent.FluentFuture; import com.google.common.util.concurrent.ListenableFuture; @@ -81,6 +80,7 @@ import java.util.concurrent.ExecutionException; import java.util.stream.Collectors; import java.util.stream.Stream; +import static com.google.common.util.concurrent.MoreExecutors.directExecutor; import static org.thingsboard.server.dao.service.Validator.validateEntityDataPageLink; import static org.thingsboard.server.dao.service.Validator.validateId; @@ -97,8 +97,8 @@ public class BaseAlarmService extends AbstractCachedEntityService>> findEntityAsync(TenantId tenantId, EntityId entityId) { + return FluentFuture.from(findAlarmByIdAsync(tenantId, new AlarmId(entityId.getId()))) + .transform(Optional::ofNullable, directExecutor()); + } + @Override public EntityType getEntityType() { return EntityType.ALARM; diff --git a/dao/src/main/java/org/thingsboard/server/dao/asset/AssetProfileServiceImpl.java b/dao/src/main/java/org/thingsboard/server/dao/asset/AssetProfileServiceImpl.java index 6bf44e4da1..4ddbabb5b6 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/asset/AssetProfileServiceImpl.java +++ b/dao/src/main/java/org/thingsboard/server/dao/asset/AssetProfileServiceImpl.java @@ -15,6 +15,7 @@ */ package org.thingsboard.server.dao.asset; +import com.google.common.util.concurrent.FluentFuture; import lombok.extern.slf4j.Slf4j; import org.hibernate.exception.ConstraintViolationException; import org.springframework.beans.factory.annotation.Autowired; @@ -49,6 +50,7 @@ import java.util.Map; import java.util.Optional; import java.util.stream.Collectors; +import static com.google.common.util.concurrent.MoreExecutors.directExecutor; import static org.thingsboard.server.dao.service.Validator.validateId; @Service("AssetProfileDaoService") @@ -323,6 +325,12 @@ public class AssetProfileServiceImpl extends CachedVersionedEntityService>> findEntityAsync(TenantId tenantId, EntityId entityId) { + return FluentFuture.from(assetProfileDao.findByIdAsync(tenantId, entityId.getId())) + .transform(Optional::ofNullable, directExecutor()); + } + @Override public EntityType getEntityType() { return EntityType.ASSET_PROFILE; diff --git a/dao/src/main/java/org/thingsboard/server/dao/asset/BaseAssetService.java b/dao/src/main/java/org/thingsboard/server/dao/asset/BaseAssetService.java index fed201d403..b8dc347a09 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/asset/BaseAssetService.java +++ b/dao/src/main/java/org/thingsboard/server/dao/asset/BaseAssetService.java @@ -15,7 +15,7 @@ */ package org.thingsboard.server.dao.asset; - +import com.google.common.util.concurrent.FluentFuture; import com.google.common.util.concurrent.Futures; import com.google.common.util.concurrent.ListenableFuture; import com.google.common.util.concurrent.MoreExecutors; @@ -62,6 +62,7 @@ import java.util.List; import java.util.Optional; import java.util.stream.Collectors; +import static com.google.common.util.concurrent.MoreExecutors.directExecutor; import static org.thingsboard.server.dao.DaoUtil.toUUIDs; import static org.thingsboard.server.dao.service.Validator.validateId; import static org.thingsboard.server.dao.service.Validator.validateIds; @@ -93,8 +94,8 @@ public class BaseAssetService extends AbstractCachedEntityService keys = new ArrayList<>(2); keys.add(new AssetCacheKey(event.getTenantId(), event.getNewName())); @@ -519,6 +520,12 @@ public class BaseAssetService extends AbstractCachedEntityService>> findEntityAsync(TenantId tenantId, EntityId entityId) { + return FluentFuture.from(findAssetByIdAsync(tenantId, new AssetId(entityId.getId()))) + .transform(Optional::ofNullable, directExecutor()); + } + @Override public long countByTenantId(TenantId tenantId) { return assetDao.countByTenantId(tenantId); diff --git a/dao/src/main/java/org/thingsboard/server/dao/cf/BaseCalculatedFieldService.java b/dao/src/main/java/org/thingsboard/server/dao/cf/BaseCalculatedFieldService.java index c0cb886747..c9a53d795e 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/cf/BaseCalculatedFieldService.java +++ b/dao/src/main/java/org/thingsboard/server/dao/cf/BaseCalculatedFieldService.java @@ -15,6 +15,7 @@ */ package org.thingsboard.server.dao.cf; +import com.google.common.util.concurrent.FluentFuture; import lombok.RequiredArgsConstructor; import lombok.extern.slf4j.Slf4j; import org.springframework.stereotype.Service; @@ -38,6 +39,7 @@ import org.thingsboard.server.dao.service.DataValidator; import java.util.List; import java.util.Optional; +import static com.google.common.util.concurrent.MoreExecutors.directExecutor; import static org.thingsboard.server.dao.service.Validator.validateId; import static org.thingsboard.server.dao.service.Validator.validatePageLink; @@ -234,6 +236,12 @@ public class BaseCalculatedFieldService extends AbstractEntityService implements return Optional.ofNullable(findById(tenantId, new CalculatedFieldId(entityId.getId()))); } + @Override + public FluentFuture>> findEntityAsync(TenantId tenantId, EntityId entityId) { + return FluentFuture.from(calculatedFieldDao.findByIdAsync(tenantId, entityId.getId())) + .transform(Optional::ofNullable, directExecutor()); + } + @Override public EntityType getEntityType() { return EntityType.CALCULATED_FIELD; diff --git a/dao/src/main/java/org/thingsboard/server/dao/customer/CustomerServiceImpl.java b/dao/src/main/java/org/thingsboard/server/dao/customer/CustomerServiceImpl.java index b06902d95e..9aee422030 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/customer/CustomerServiceImpl.java +++ b/dao/src/main/java/org/thingsboard/server/dao/customer/CustomerServiceImpl.java @@ -16,6 +16,7 @@ package org.thingsboard.server.dao.customer; import com.fasterxml.jackson.databind.JsonNode; +import com.google.common.util.concurrent.FluentFuture; import com.google.common.util.concurrent.ListenableFuture; import lombok.extern.slf4j.Slf4j; import org.springframework.beans.factory.annotation.Autowired; @@ -55,6 +56,7 @@ import java.util.ArrayList; import java.util.List; import java.util.Optional; +import static com.google.common.util.concurrent.MoreExecutors.directExecutor; import static org.thingsboard.server.dao.service.Validator.validateId; @Service("CustomerDaoService") @@ -285,6 +287,12 @@ public class CustomerServiceImpl extends AbstractCachedEntityService>> findEntityAsync(TenantId tenantId, EntityId entityId) { + return FluentFuture.from(findCustomerByIdAsync(tenantId, new CustomerId(entityId.getId()))) + .transform(Optional::ofNullable, directExecutor()); + } + @Override public long countByTenantId(TenantId tenantId) { return customerDao.countByTenantId(tenantId); diff --git a/dao/src/main/java/org/thingsboard/server/dao/dashboard/DashboardServiceImpl.java b/dao/src/main/java/org/thingsboard/server/dao/dashboard/DashboardServiceImpl.java index 1b21310237..9d8abe467c 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/dashboard/DashboardServiceImpl.java +++ b/dao/src/main/java/org/thingsboard/server/dao/dashboard/DashboardServiceImpl.java @@ -15,6 +15,7 @@ */ package org.thingsboard.server.dao.dashboard; +import com.google.common.util.concurrent.FluentFuture; import com.google.common.util.concurrent.ListenableFuture; import lombok.RequiredArgsConstructor; import lombok.extern.slf4j.Slf4j; @@ -61,6 +62,7 @@ import java.util.List; import java.util.Map; import java.util.Optional; +import static com.google.common.util.concurrent.MoreExecutors.directExecutor; import static org.thingsboard.server.dao.service.Validator.validateId; @Service("DashboardDaoService") @@ -70,6 +72,7 @@ public class DashboardServiceImpl extends AbstractEntityService implements Dashb public static final String INCORRECT_DASHBOARD_ID = "Incorrect dashboardId "; public static final String INCORRECT_TENANT_ID = "Incorrect tenantId "; + @Autowired private DashboardDao dashboardDao; @@ -424,6 +427,12 @@ public class DashboardServiceImpl extends AbstractEntityService implements Dashb return Optional.ofNullable(findDashboardById(tenantId, new DashboardId(entityId.getId()))); } + @Override + public FluentFuture>> findEntityAsync(TenantId tenantId, EntityId entityId) { + return FluentFuture.from(findDashboardByIdAsync(tenantId, new DashboardId(entityId.getId()))) + .transform(Optional::ofNullable, directExecutor()); + } + @Override public long countByTenantId(TenantId tenantId) { return dashboardDao.countByTenantId(tenantId); diff --git a/dao/src/main/java/org/thingsboard/server/dao/device/DeviceProfileServiceImpl.java b/dao/src/main/java/org/thingsboard/server/dao/device/DeviceProfileServiceImpl.java index 2ebcf046a6..f782585888 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/device/DeviceProfileServiceImpl.java +++ b/dao/src/main/java/org/thingsboard/server/dao/device/DeviceProfileServiceImpl.java @@ -15,6 +15,7 @@ */ package org.thingsboard.server.dao.device; +import com.google.common.util.concurrent.FluentFuture; import lombok.RequiredArgsConstructor; import lombok.extern.slf4j.Slf4j; import org.springframework.beans.factory.annotation.Autowired; @@ -64,6 +65,7 @@ import java.util.regex.Matcher; import java.util.regex.Pattern; import java.util.stream.Collectors; +import static com.google.common.util.concurrent.MoreExecutors.directExecutor; import static org.thingsboard.server.dao.service.Validator.validateId; import static org.thingsboard.server.dao.service.Validator.validateString; @@ -369,6 +371,12 @@ public class DeviceProfileServiceImpl extends CachedVersionedEntityService>> findEntityAsync(TenantId tenantId, EntityId entityId) { + return FluentFuture.from(deviceProfileDao.findByIdAsync(tenantId, entityId.getId())) + .transform(Optional::ofNullable, directExecutor()); + } + @Override public EntityType getEntityType() { return EntityType.DEVICE_PROFILE; @@ -432,8 +440,7 @@ public class DeviceProfileServiceImpl extends CachedVersionedEntityService 1) { return EncryptionUtil.certTrimNewLinesForChainInDeviceProfile(certificateValue); } - } catch (CertificateException ignored) { - } + } catch (CertificateException ignored) {} return EncryptionUtil.certTrimNewLines(certificateValue); } diff --git a/dao/src/main/java/org/thingsboard/server/dao/device/DeviceServiceImpl.java b/dao/src/main/java/org/thingsboard/server/dao/device/DeviceServiceImpl.java index 6d993f3e3d..4a6e6ce9be 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/device/DeviceServiceImpl.java +++ b/dao/src/main/java/org/thingsboard/server/dao/device/DeviceServiceImpl.java @@ -16,6 +16,7 @@ package org.thingsboard.server.dao.device; import com.fasterxml.jackson.databind.node.ObjectNode; +import com.google.common.util.concurrent.FluentFuture; import com.google.common.util.concurrent.Futures; import com.google.common.util.concurrent.ListenableFuture; import com.google.common.util.concurrent.MoreExecutors; @@ -89,6 +90,7 @@ import java.util.List; import java.util.Optional; import java.util.UUID; +import static com.google.common.util.concurrent.MoreExecutors.directExecutor; import static org.thingsboard.server.dao.DaoUtil.toUUIDs; import static org.thingsboard.server.dao.service.Validator.validateId; import static org.thingsboard.server.dao.service.Validator.validateIds; @@ -126,23 +128,21 @@ public class DeviceServiceImpl extends CachedVersionedEntityService INCORRECT_DEVICE_ID + id); - if (TenantId.SYS_TENANT_ID.equals(tenantId)) { - return cache.get(new DeviceCacheKey(deviceId), - () -> deviceDao.findById(tenantId, deviceId.getId())); - } else { - return cache.get(new DeviceCacheKey(tenantId, deviceId), - () -> deviceDao.findDeviceByTenantIdAndId(tenantId, deviceId.getId())); - } + return findDeviceByIdInternal(tenantId, deviceId); } @Override public ListenableFuture findDeviceByIdAsync(TenantId tenantId, DeviceId deviceId) { log.trace("Executing findDeviceByIdAsync [{}]", deviceId); validateId(deviceId, id -> INCORRECT_DEVICE_ID + id); + return executor.submit(() -> findDeviceByIdInternal(tenantId, deviceId)); + } + + private Device findDeviceByIdInternal(TenantId tenantId, DeviceId deviceId) { if (TenantId.SYS_TENANT_ID.equals(tenantId)) { - return deviceDao.findByIdAsync(tenantId, deviceId.getId()); + return cache.get(new DeviceCacheKey(deviceId), () -> deviceDao.findById(tenantId, deviceId.getId())); } else { - return deviceDao.findDeviceByTenantIdAndIdAsync(tenantId, deviceId.getId()); + return cache.get(new DeviceCacheKey(tenantId, deviceId), () -> deviceDao.findDeviceByTenantIdAndId(tenantId, deviceId.getId())); } } @@ -256,8 +256,8 @@ public class DeviceServiceImpl extends CachedVersionedEntityService toEvict = new ArrayList<>(3); toEvict.add(new DeviceCacheKey(event.getTenantId(), event.getNewName())); @@ -729,6 +729,12 @@ public class DeviceServiceImpl extends CachedVersionedEntityService>> findEntityAsync(TenantId tenantId, EntityId entityId) { + return FluentFuture.from(findDeviceByIdAsync(tenantId, new DeviceId(entityId.getId()))) + .transform(Optional::ofNullable, directExecutor()); + } + @Override public EntityType getEntityType() { return EntityType.DEVICE; diff --git a/dao/src/main/java/org/thingsboard/server/dao/domain/DomainServiceImpl.java b/dao/src/main/java/org/thingsboard/server/dao/domain/DomainServiceImpl.java index c12d4d915e..e4e569c00a 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/domain/DomainServiceImpl.java +++ b/dao/src/main/java/org/thingsboard/server/dao/domain/DomainServiceImpl.java @@ -15,6 +15,7 @@ */ package org.thingsboard.server.dao.domain; +import com.google.common.util.concurrent.FluentFuture; import lombok.extern.slf4j.Slf4j; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.stereotype.Service; @@ -39,17 +40,16 @@ import org.thingsboard.server.dao.service.validator.DomainDataValidator; import java.util.Comparator; import java.util.List; -import java.util.Map; import java.util.Optional; import java.util.Set; import java.util.stream.Collectors; +import static com.google.common.util.concurrent.MoreExecutors.directExecutor; + @Slf4j @Service public class DomainServiceImpl extends AbstractEntityService implements DomainService { - public static final String INCORRECT_TENANT_ID = "Incorrect tenantId "; - @Autowired private OAuth2ClientDao oauth2ClientDao; @Autowired @@ -66,8 +66,7 @@ public class DomainServiceImpl extends AbstractEntityService implements DomainSe eventPublisher.publishEvent(SaveEntityEvent.builder().tenantId(tenantId).entityId(savedDomain.getId()).entity(savedDomain).build()); return savedDomain; } catch (Exception e) { - checkConstraintViolation(e, - Map.of("domain_name_key", "Domain with such name and scheme already exists!")); + checkConstraintViolation(e, "domain_name_key", "Domain with such name and scheme already exists!"); throw e; } } @@ -142,6 +141,12 @@ public class DomainServiceImpl extends AbstractEntityService implements DomainSe return Optional.ofNullable(findDomainById(tenantId, new DomainId(entityId.getId()))); } + @Override + public FluentFuture>> findEntityAsync(TenantId tenantId, EntityId entityId) { + return FluentFuture.from(domainDao.findByIdAsync(tenantId, entityId.getId())) + .transform(Optional::ofNullable, directExecutor()); + } + @Override @Transactional public void deleteEntity(TenantId tenantId, EntityId id, boolean force) { @@ -163,4 +168,5 @@ public class DomainServiceImpl extends AbstractEntityService implements DomainSe public EntityType getEntityType() { return EntityType.DOMAIN; } + } diff --git a/dao/src/main/java/org/thingsboard/server/dao/edge/EdgeServiceImpl.java b/dao/src/main/java/org/thingsboard/server/dao/edge/EdgeServiceImpl.java index 0655d05572..f0658ba658 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/edge/EdgeServiceImpl.java +++ b/dao/src/main/java/org/thingsboard/server/dao/edge/EdgeServiceImpl.java @@ -17,6 +17,7 @@ package org.thingsboard.server.dao.edge; import com.fasterxml.jackson.databind.node.ArrayNode; import com.fasterxml.jackson.databind.node.ObjectNode; +import com.google.common.util.concurrent.FluentFuture; import com.google.common.util.concurrent.Futures; import com.google.common.util.concurrent.ListenableFuture; import com.google.common.util.concurrent.MoreExecutors; @@ -83,6 +84,7 @@ import java.util.Optional; import java.util.UUID; import java.util.stream.Collectors; +import static com.google.common.util.concurrent.MoreExecutors.directExecutor; import static org.thingsboard.server.dao.DaoUtil.toUUIDs; import static org.thingsboard.server.dao.edge.BaseRelatedEdgesService.RELATED_EDGES_CACHE_ITEMS; import static org.thingsboard.server.dao.service.Validator.validateId; @@ -137,8 +139,8 @@ public class EdgeServiceImpl extends AbstractCachedEntityService keys = new ArrayList<>(2); keys.add(new EdgeCacheKey(event.getTenantId(), event.getNewName())); @@ -629,6 +631,12 @@ public class EdgeServiceImpl extends AbstractCachedEntityService>> findEntityAsync(TenantId tenantId, EntityId entityId) { + return FluentFuture.from(findEdgeByIdAsync(tenantId, new EdgeId(entityId.getId()))) + .transform(Optional::ofNullable, directExecutor()); + } + @Override public EntityType getEntityType() { return EntityType.EDGE; diff --git a/dao/src/main/java/org/thingsboard/server/dao/entity/BaseEntityService.java b/dao/src/main/java/org/thingsboard/server/dao/entity/BaseEntityService.java index 2fca546fc9..b512f353f5 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/entity/BaseEntityService.java +++ b/dao/src/main/java/org/thingsboard/server/dao/entity/BaseEntityService.java @@ -15,6 +15,7 @@ */ package org.thingsboard.server.dao.entity; +import com.google.common.util.concurrent.FluentFuture; import lombok.extern.slf4j.Slf4j; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.context.annotation.Lazy; @@ -67,15 +68,13 @@ import java.util.concurrent.ExecutionException; import java.util.function.Function; import java.util.stream.Collectors; +import static com.google.common.util.concurrent.MoreExecutors.directExecutor; import static org.thingsboard.server.common.data.id.EntityId.NULL_UUID; import static org.thingsboard.server.common.data.query.EntityFilterType.ENTITY_NAME; import static org.thingsboard.server.common.data.query.EntityFilterType.ENTITY_TYPE; import static org.thingsboard.server.dao.service.Validator.validateEntityDataPageLink; import static org.thingsboard.server.dao.service.Validator.validateId; -/** - * Created by ashvayka on 04.05.17. - */ @Service @Slf4j public class BaseEntityService extends AbstractEntityService implements EntityService { @@ -93,7 +92,7 @@ public class BaseEntityService extends AbstractEntityService implements EntitySe @Autowired @Lazy - EntityServiceRegistry entityServiceRegistry; + private EntityServiceRegistry entityServiceRegistry; @Autowired private EdqsService edqsService; @@ -198,6 +197,11 @@ public class BaseEntityService extends AbstractEntityService implements EntitySe return fetchAndConvert(tenantId, entityId, this::getCustomerId); } + @Override + public FluentFuture> fetchEntityCustomerIdAsync(TenantId tenantId, EntityId entityId) { + return fetchAndConvertAsync(tenantId, entityId, this::getCustomerId); + } + @Override public Optional fetchNameLabelAndCustomerDetails(TenantId tenantId, EntityId entityId) { log.trace("Executing fetchNameLabelAndCustomerDetails [{}]", entityId); @@ -239,6 +243,12 @@ public class BaseEntityService extends AbstractEntityService implements EntitySe return entityOpt.map(converter); } + private FluentFuture> fetchAndConvertAsync(TenantId tenantId, EntityId entityId, Function, T> converter) { + EntityDaoService entityDaoService = entityServiceRegistry.getServiceByEntityType(entityId.getEntityType()); + return entityDaoService.findEntityAsync(tenantId, entityId) + .transform(entityOpt -> entityOpt.map(converter), directExecutor()); + } + private String getName(HasId entity) { return entity instanceof HasName ? ((HasName) entity).getName() : null; } @@ -329,7 +339,7 @@ public class BaseEntityService extends AbstractEntityService implements EntitySe } if ((query.getEntityFields() == null || query.getEntityFields().isEmpty()) && - (query.getLatestValues() == null || query.getLatestValues().isEmpty())) { + (query.getLatestValues() == null || query.getLatestValues().isEmpty())) { return false; } diff --git a/dao/src/main/java/org/thingsboard/server/dao/entityview/EntityViewServiceImpl.java b/dao/src/main/java/org/thingsboard/server/dao/entityview/EntityViewServiceImpl.java index 0e742e20db..3128036e0e 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/entityview/EntityViewServiceImpl.java +++ b/dao/src/main/java/org/thingsboard/server/dao/entityview/EntityViewServiceImpl.java @@ -16,6 +16,7 @@ package org.thingsboard.server.dao.entityview; import com.google.common.base.Function; +import com.google.common.util.concurrent.FluentFuture; import com.google.common.util.concurrent.Futures; import com.google.common.util.concurrent.ListenableFuture; import com.google.common.util.concurrent.MoreExecutors; @@ -60,13 +61,11 @@ import java.util.List; import java.util.Optional; import java.util.stream.Collectors; +import static com.google.common.util.concurrent.MoreExecutors.directExecutor; import static org.thingsboard.server.dao.service.Validator.validateId; import static org.thingsboard.server.dao.service.Validator.validatePageLink; import static org.thingsboard.server.dao.service.Validator.validateString; -/** - * Created by Victor Basanets on 8/28/2017. - */ @Service("EntityViewDaoService") @Slf4j public class EntityViewServiceImpl extends CachedVersionedEntityService implements EntityViewService { @@ -176,6 +175,17 @@ public class EntityViewServiceImpl extends CachedVersionedEntityService INCORRECT_ENTITY_VIEW_ID + id); + return findEntityViewByIdInternal(tenantId, entityViewId, putInCache); + } + + @Override + public ListenableFuture findEntityViewByIdAsync(TenantId tenantId, EntityViewId entityViewId) { + log.trace("Executing findEntityViewByIdAsync [{}]", entityViewId); + validateId(entityViewId, id -> INCORRECT_ENTITY_VIEW_ID + id); + return service.submit(() -> findEntityViewByIdInternal(tenantId, entityViewId, true)); + } + + private EntityView findEntityViewByIdInternal(TenantId tenantId, EntityViewId entityViewId, boolean putInCache) { EntityViewCacheValue value = cache.get(EntityViewCacheKey.byId(entityViewId), () -> { EntityView entityView = entityViewDao.findById(tenantId, entityViewId.getId()); return new EntityViewCacheValue(entityView, null); @@ -190,7 +200,6 @@ public class EntityViewServiceImpl extends CachedVersionedEntityService entityViewDao.findEntityViewByTenantIdAndName(tenantId.getId(), name).orElse(null) , EntityViewCacheValue::getEntityView, v -> new EntityViewCacheValue(v, null), true); - } @Override @@ -307,13 +316,6 @@ public class EntityViewServiceImpl extends CachedVersionedEntityService findEntityViewByIdAsync(TenantId tenantId, EntityViewId entityViewId) { - log.trace("Executing findEntityViewByIdAsync [{}]", entityViewId); - validateId(entityViewId, id -> INCORRECT_ENTITY_VIEW_ID + id); - return entityViewDao.findByIdAsync(tenantId, entityViewId.getId()); - } - @Override public ListenableFuture> findEntityViewsByTenantIdAndEntityIdAsync(TenantId tenantId, EntityId entityId) { log.trace("Executing findEntityViewsByTenantIdAndEntityIdAsync, tenantId [{}], entityId [{}]", tenantId, entityId); @@ -479,6 +481,12 @@ public class EntityViewServiceImpl extends CachedVersionedEntityService>> findEntityAsync(TenantId tenantId, EntityId entityId) { + return FluentFuture.from(findEntityViewByIdAsync(tenantId, new EntityViewId(entityId.getId()))) + .transform(Optional::ofNullable, directExecutor()); + } + @Override public EntityType getEntityType() { return EntityType.ENTITY_VIEW; diff --git a/dao/src/main/java/org/thingsboard/server/dao/job/DefaultJobService.java b/dao/src/main/java/org/thingsboard/server/dao/job/DefaultJobService.java index 360aa0063b..b435066cd9 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/job/DefaultJobService.java +++ b/dao/src/main/java/org/thingsboard/server/dao/job/DefaultJobService.java @@ -15,6 +15,7 @@ */ package org.thingsboard.server.dao.job; +import com.google.common.util.concurrent.FluentFuture; import lombok.RequiredArgsConstructor; import lombok.extern.slf4j.Slf4j; import org.springframework.stereotype.Service; @@ -44,6 +45,7 @@ import java.util.Optional; import java.util.Set; import java.util.stream.Collectors; +import static com.google.common.util.concurrent.MoreExecutors.directExecutor; import static org.thingsboard.server.common.data.job.JobStatus.CANCELLED; import static org.thingsboard.server.common.data.job.JobStatus.COMPLETED; import static org.thingsboard.server.common.data.job.JobStatus.FAILED; @@ -245,6 +247,12 @@ public class DefaultJobService extends AbstractEntityService implements JobServi return Optional.ofNullable(findJobById(tenantId, (JobId) entityId)); } + @Override + public FluentFuture>> findEntityAsync(TenantId tenantId, EntityId entityId) { + return FluentFuture.from(jobDao.findByIdAsync(tenantId, entityId.getId())) + .transform(Optional::ofNullable, directExecutor()); + } + @Override public void deleteEntity(TenantId tenantId, EntityId id, boolean force) { jobDao.removeById(tenantId, id.getId()); diff --git a/dao/src/main/java/org/thingsboard/server/dao/mobile/MobileAppBundleServiceImpl.java b/dao/src/main/java/org/thingsboard/server/dao/mobile/MobileAppBundleServiceImpl.java index d20623ad66..98f4a612e7 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/mobile/MobileAppBundleServiceImpl.java +++ b/dao/src/main/java/org/thingsboard/server/dao/mobile/MobileAppBundleServiceImpl.java @@ -15,6 +15,7 @@ */ package org.thingsboard.server.dao.mobile; +import com.google.common.util.concurrent.FluentFuture; import lombok.extern.slf4j.Slf4j; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.stereotype.Service; @@ -40,11 +41,11 @@ import org.thingsboard.server.dao.service.DataValidator; import java.util.Comparator; import java.util.List; -import java.util.Map; import java.util.Optional; import java.util.Set; import java.util.stream.Collectors; +import static com.google.common.util.concurrent.MoreExecutors.directExecutor; import static org.thingsboard.server.dao.service.Validator.checkNotNull; @Slf4j @@ -60,7 +61,6 @@ public class MobileAppBundleServiceImpl extends AbstractEntityService implements @Autowired private DataValidator mobileAppBundleDataValidator; - @Override public MobileAppBundle saveMobileAppBundle(TenantId tenantId, MobileAppBundle mobileAppBundle) { log.trace("Executing saveMobileAppBundle [{}]", mobileAppBundle); @@ -71,8 +71,8 @@ public class MobileAppBundleServiceImpl extends AbstractEntityService implements return savedMobileApp; } catch (Exception e) { checkConstraintViolation(e, - Map.of("mobile_app_bundle_android_app_id_key", "Android mobile app is already configured in another bundle!", - "mobile_app_bundle_ios_app_id_key", "IOS mobile app is already configured in another bundle!")); + "mobile_app_bundle_android_app_id_key", "Android mobile app is already configured in another bundle!", + "mobile_app_bundle_ios_app_id_key", "IOS mobile app is already configured in another bundle!"); throw e; } } @@ -149,6 +149,12 @@ public class MobileAppBundleServiceImpl extends AbstractEntityService implements return Optional.ofNullable(findMobileAppBundleById(tenantId, new MobileAppBundleId(entityId.getId()))); } + @Override + public FluentFuture>> findEntityAsync(TenantId tenantId, EntityId entityId) { + return FluentFuture.from(mobileAppBundleDao.findByIdAsync(tenantId, entityId.getId())) + .transform(Optional::ofNullable, directExecutor()); + } + @Override @Transactional public void deleteEntity(TenantId tenantId, EntityId id, boolean force) { @@ -167,4 +173,5 @@ public class MobileAppBundleServiceImpl extends AbstractEntityService implements .collect(Collectors.toList()); mobileAppBundleInfo.setOauth2ClientInfos(clients); } + } diff --git a/dao/src/main/java/org/thingsboard/server/dao/mobile/MobileAppServiceImpl.java b/dao/src/main/java/org/thingsboard/server/dao/mobile/MobileAppServiceImpl.java index f4e803b28e..b7785496d4 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/mobile/MobileAppServiceImpl.java +++ b/dao/src/main/java/org/thingsboard/server/dao/mobile/MobileAppServiceImpl.java @@ -15,6 +15,7 @@ */ package org.thingsboard.server.dao.mobile; +import com.google.common.util.concurrent.FluentFuture; import lombok.extern.slf4j.Slf4j; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.stereotype.Service; @@ -35,9 +36,10 @@ import org.thingsboard.server.dao.eventsourcing.SaveEntityEvent; import org.thingsboard.server.dao.service.DataValidator; import org.thingsboard.server.dao.service.Validator; -import java.util.Map; import java.util.Optional; +import static com.google.common.util.concurrent.MoreExecutors.directExecutor; + @Slf4j @Service public class MobileAppServiceImpl extends AbstractEntityService implements MobileAppService { @@ -58,8 +60,7 @@ public class MobileAppServiceImpl extends AbstractEntityService implements Mobil eventPublisher.publishEvent(SaveEntityEvent.builder().tenantId(tenantId).entity(savedMobileApp).build()); return savedMobileApp; } catch (Exception e) { - checkConstraintViolation(e, - Map.of("mobile_app_pkg_name_platform_unq_key", "Mobile app with such package name and platform already exists!")); + checkConstraintViolation(e, "mobile_app_pkg_name_platform_unq_key", "Mobile app with such package name and platform already exists!"); throw e; } } @@ -88,6 +89,12 @@ public class MobileAppServiceImpl extends AbstractEntityService implements Mobil return Optional.ofNullable(findMobileAppById(tenantId, new MobileAppId(entityId.getId()))); } + @Override + public FluentFuture>> findEntityAsync(TenantId tenantId, EntityId entityId) { + return FluentFuture.from(mobileAppDao.findByIdAsync(tenantId, entityId.getId())) + .transform(Optional::ofNullable, directExecutor()); + } + @Override @Transactional public void deleteEntity(TenantId tenantId, EntityId id, boolean force) { @@ -117,4 +124,5 @@ public class MobileAppServiceImpl extends AbstractEntityService implements Mobil public EntityType getEntityType() { return EntityType.MOBILE_APP; } + } diff --git a/dao/src/main/java/org/thingsboard/server/dao/notification/DefaultNotificationRequestService.java b/dao/src/main/java/org/thingsboard/server/dao/notification/DefaultNotificationRequestService.java index ffed043fe2..35c7839429 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/notification/DefaultNotificationRequestService.java +++ b/dao/src/main/java/org/thingsboard/server/dao/notification/DefaultNotificationRequestService.java @@ -15,6 +15,7 @@ */ package org.thingsboard.server.dao.notification; +import com.google.common.util.concurrent.FluentFuture; import lombok.RequiredArgsConstructor; import lombok.extern.slf4j.Slf4j; import org.springframework.context.ApplicationEventPublisher; @@ -38,6 +39,8 @@ import org.thingsboard.server.dao.service.DataValidator; import java.util.List; import java.util.Optional; +import static com.google.common.util.concurrent.MoreExecutors.directExecutor; + @Service @Slf4j @RequiredArgsConstructor @@ -129,13 +132,17 @@ public class DefaultNotificationRequestService implements NotificationRequestSer return Optional.ofNullable(findNotificationRequestById(tenantId, new NotificationRequestId(entityId.getId()))); } + @Override + public FluentFuture>> findEntityAsync(TenantId tenantId, EntityId entityId) { + return FluentFuture.from(notificationRequestDao.findByIdAsync(tenantId, entityId.getId())) + .transform(Optional::ofNullable, directExecutor()); + } + @Override public EntityType getEntityType() { return EntityType.NOTIFICATION_REQUEST; } - private static class NotificationRequestValidator extends DataValidator { - - } + private static class NotificationRequestValidator extends DataValidator {} } diff --git a/dao/src/main/java/org/thingsboard/server/dao/notification/DefaultNotificationRuleService.java b/dao/src/main/java/org/thingsboard/server/dao/notification/DefaultNotificationRuleService.java index c349b84f0a..b2f79c9555 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/notification/DefaultNotificationRuleService.java +++ b/dao/src/main/java/org/thingsboard/server/dao/notification/DefaultNotificationRuleService.java @@ -15,6 +15,7 @@ */ package org.thingsboard.server.dao.notification; +import com.google.common.util.concurrent.FluentFuture; import lombok.RequiredArgsConstructor; import org.springframework.stereotype.Service; import org.thingsboard.server.common.data.EntityType; @@ -33,9 +34,10 @@ import org.thingsboard.server.dao.eventsourcing.DeleteEntityEvent; import org.thingsboard.server.dao.eventsourcing.SaveEntityEvent; import java.util.List; -import java.util.Map; import java.util.Optional; +import static com.google.common.util.concurrent.MoreExecutors.directExecutor; + @Service @RequiredArgsConstructor public class DefaultNotificationRuleService extends AbstractEntityService implements NotificationRuleService, EntityDaoService { @@ -56,9 +58,7 @@ public class DefaultNotificationRuleService extends AbstractEntityService implem .created(notificationRule.getId() == null).build()); return savedRule; } catch (Exception e) { - checkConstraintViolation(e, Map.of( - "uq_notification_rule_name", "Notification rule with such name already exists" - )); + checkConstraintViolation(e, "uq_notification_rule_name", "Notification rule with such name already exists"); throw e; } } @@ -114,6 +114,12 @@ public class DefaultNotificationRuleService extends AbstractEntityService implem return Optional.ofNullable(findNotificationRuleById(tenantId, new NotificationRuleId(entityId.getId()))); } + @Override + public FluentFuture>> findEntityAsync(TenantId tenantId, EntityId entityId) { + return FluentFuture.from(notificationRuleDao.findByIdAsync(tenantId, entityId.getId())) + .transform(Optional::ofNullable, directExecutor()); + } + @Override public EntityType getEntityType() { return EntityType.NOTIFICATION_RULE; diff --git a/dao/src/main/java/org/thingsboard/server/dao/notification/DefaultNotificationService.java b/dao/src/main/java/org/thingsboard/server/dao/notification/DefaultNotificationService.java index 99efd18aee..8c5a051771 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/notification/DefaultNotificationService.java +++ b/dao/src/main/java/org/thingsboard/server/dao/notification/DefaultNotificationService.java @@ -15,6 +15,7 @@ */ package org.thingsboard.server.dao.notification; +import com.google.common.util.concurrent.FluentFuture; import lombok.RequiredArgsConstructor; import lombok.extern.slf4j.Slf4j; import org.springframework.stereotype.Service; @@ -37,6 +38,8 @@ import org.thingsboard.server.dao.sql.query.EntityKeyMapping; import java.util.Optional; import java.util.Set; +import static com.google.common.util.concurrent.MoreExecutors.directExecutor; + @Service @Slf4j @RequiredArgsConstructor @@ -95,6 +98,12 @@ public class DefaultNotificationService implements NotificationService, EntityDa return Optional.ofNullable(findNotificationById(tenantId, new NotificationId(entityId.getId()))); } + @Override + public FluentFuture>> findEntityAsync(TenantId tenantId, EntityId entityId) { + return FluentFuture.from(notificationDao.findByIdAsync(tenantId, entityId.getId())) + .transform(Optional::ofNullable, directExecutor()); + } + @Override public EntityType getEntityType() { return EntityType.NOTIFICATION; diff --git a/dao/src/main/java/org/thingsboard/server/dao/notification/DefaultNotificationTargetService.java b/dao/src/main/java/org/thingsboard/server/dao/notification/DefaultNotificationTargetService.java index e873b18c76..02a7e4b891 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/notification/DefaultNotificationTargetService.java +++ b/dao/src/main/java/org/thingsboard/server/dao/notification/DefaultNotificationTargetService.java @@ -15,6 +15,7 @@ */ package org.thingsboard.server.dao.notification; +import com.google.common.util.concurrent.FluentFuture; import lombok.RequiredArgsConstructor; import lombok.extern.slf4j.Slf4j; import org.springframework.stereotype.Service; @@ -47,13 +48,14 @@ import org.thingsboard.server.dao.eventsourcing.SaveEntityEvent; import org.thingsboard.server.dao.user.UserService; import java.util.List; -import java.util.Map; import java.util.Objects; import java.util.Optional; import java.util.stream.Collectors; import static org.apache.commons.collections4.CollectionUtils.isNotEmpty; +import static com.google.common.util.concurrent.MoreExecutors.directExecutor; + @Service @Slf4j @RequiredArgsConstructor @@ -72,9 +74,7 @@ public class DefaultNotificationTargetService extends AbstractEntityService impl .created(notificationTarget.getId() == null).build()); return savedTarget; } catch (Exception e) { - checkConstraintViolation(e, Map.of( - "uq_notification_target_name", "Recipients group with such name already exists" - )); + checkConstraintViolation(e, "uq_notification_target_name", "Recipients group with such name already exists"); throw e; } } @@ -229,6 +229,12 @@ public class DefaultNotificationTargetService extends AbstractEntityService impl return Optional.ofNullable(findNotificationTargetById(tenantId, new NotificationTargetId(entityId.getId()))); } + @Override + public FluentFuture>> findEntityAsync(TenantId tenantId, EntityId entityId) { + return FluentFuture.from(notificationTargetDao.findByIdAsync(tenantId, entityId.getId())) + .transform(Optional::ofNullable, directExecutor()); + } + @Override public EntityType getEntityType() { return EntityType.NOTIFICATION_TARGET; diff --git a/dao/src/main/java/org/thingsboard/server/dao/notification/DefaultNotificationTemplateService.java b/dao/src/main/java/org/thingsboard/server/dao/notification/DefaultNotificationTemplateService.java index 8e98714b85..4218e06881 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/notification/DefaultNotificationTemplateService.java +++ b/dao/src/main/java/org/thingsboard/server/dao/notification/DefaultNotificationTemplateService.java @@ -15,6 +15,7 @@ */ package org.thingsboard.server.dao.notification; +import com.google.common.util.concurrent.FluentFuture; import lombok.RequiredArgsConstructor; import org.springframework.stereotype.Service; import org.thingsboard.server.common.data.EntityType; @@ -37,6 +38,8 @@ import java.util.List; import java.util.Map; import java.util.Optional; +import static com.google.common.util.concurrent.MoreExecutors.directExecutor; + @Service @RequiredArgsConstructor public class DefaultNotificationTemplateService extends AbstractEntityService implements NotificationTemplateService, EntityDaoService { @@ -144,6 +147,12 @@ public class DefaultNotificationTemplateService extends AbstractEntityService im return Optional.ofNullable(findNotificationTemplateById(tenantId, new NotificationTemplateId(entityId.getId()))); } + @Override + public FluentFuture>> findEntityAsync(TenantId tenantId, EntityId entityId) { + return FluentFuture.from(notificationTemplateDao.findByIdAsync(tenantId, entityId.getId())) + .transform(Optional::ofNullable, directExecutor()); + } + @Override public EntityType getEntityType() { return EntityType.NOTIFICATION_TEMPLATE; diff --git a/dao/src/main/java/org/thingsboard/server/dao/oauth2/OAuth2ClientServiceImpl.java b/dao/src/main/java/org/thingsboard/server/dao/oauth2/OAuth2ClientServiceImpl.java index 861633ec3b..a373d36580 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/oauth2/OAuth2ClientServiceImpl.java +++ b/dao/src/main/java/org/thingsboard/server/dao/oauth2/OAuth2ClientServiceImpl.java @@ -15,6 +15,7 @@ */ package org.thingsboard.server.dao.oauth2; +import com.google.common.util.concurrent.FluentFuture; import jakarta.transaction.Transactional; import lombok.extern.slf4j.Slf4j; import org.springframework.beans.factory.annotation.Autowired; @@ -41,6 +42,7 @@ import java.util.List; import java.util.Optional; import java.util.stream.Collectors; +import static com.google.common.util.concurrent.MoreExecutors.directExecutor; @Slf4j @Service("OAuth2ClientService") @@ -107,7 +109,6 @@ public class OAuth2ClientServiceImpl extends AbstractEntityService implements OA .tenantId(tenantId) .entityId(oAuth2ClientId) .build()); - } @Override @@ -149,6 +150,12 @@ public class OAuth2ClientServiceImpl extends AbstractEntityService implements OA return Optional.ofNullable(findOAuth2ClientById(tenantId, new OAuth2ClientId(entityId.getId()))); } + @Override + public FluentFuture>> findEntityAsync(TenantId tenantId, EntityId entityId) { + return FluentFuture.from(oauth2ClientDao.findByIdAsync(tenantId, entityId.getId())) + .transform(Optional::ofNullable, directExecutor()); + } + @Override @Transactional public void deleteEntity(TenantId tenantId, EntityId id, boolean force) { diff --git a/dao/src/main/java/org/thingsboard/server/dao/ota/BaseOtaPackageService.java b/dao/src/main/java/org/thingsboard/server/dao/ota/BaseOtaPackageService.java index 343a2485ce..16894fe56c 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/ota/BaseOtaPackageService.java +++ b/dao/src/main/java/org/thingsboard/server/dao/ota/BaseOtaPackageService.java @@ -17,6 +17,7 @@ package org.thingsboard.server.dao.ota; import com.google.common.hash.HashFunction; import com.google.common.hash.Hashing; +import com.google.common.util.concurrent.FluentFuture; import com.google.common.util.concurrent.ListenableFuture; import lombok.RequiredArgsConstructor; import lombok.extern.slf4j.Slf4j; @@ -47,6 +48,7 @@ import org.thingsboard.server.dao.service.PaginatedRemover; import java.nio.ByteBuffer; import java.util.Optional; +import static com.google.common.util.concurrent.MoreExecutors.directExecutor; import static org.thingsboard.server.dao.service.Validator.validateId; import static org.thingsboard.server.dao.service.Validator.validatePageLink; @@ -254,6 +256,12 @@ public class BaseOtaPackageService extends AbstractCachedEntityService>> findEntityAsync(TenantId tenantId, EntityId entityId) { + return FluentFuture.from(findOtaPackageInfoByIdAsync(tenantId, new OtaPackageId(entityId.getId()))) + .transform(Optional::ofNullable, directExecutor()); + } + @Override public EntityType getEntityType() { return EntityType.OTA_PACKAGE; diff --git a/dao/src/main/java/org/thingsboard/server/dao/queue/BaseQueueService.java b/dao/src/main/java/org/thingsboard/server/dao/queue/BaseQueueService.java index d36d5f2a15..248326a722 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/queue/BaseQueueService.java +++ b/dao/src/main/java/org/thingsboard/server/dao/queue/BaseQueueService.java @@ -15,6 +15,7 @@ */ package org.thingsboard.server.dao.queue; +import com.google.common.util.concurrent.FluentFuture; import lombok.RequiredArgsConstructor; import lombok.extern.slf4j.Slf4j; import org.hibernate.exception.ConstraintViolationException; @@ -42,6 +43,9 @@ import org.thingsboard.server.dao.tenant.TbTenantProfileCache; import java.util.List; import java.util.Optional; +import static com.google.common.util.concurrent.MoreExecutors.directExecutor; +import static org.thingsboard.server.dao.service.Validator.validateId; + @Service("QueueDaoService") @Slf4j @RequiredArgsConstructor @@ -127,7 +131,7 @@ public class BaseQueueService extends AbstractEntityService implements QueueServ @Override public void deleteQueuesByTenantId(TenantId tenantId) { - Validator.validateId(tenantId, "Incorrect tenant id for delete queues request."); + validateId(tenantId, __ -> "Incorrect tenant id for delete queues request."); tenantQueuesRemover.removeEntities(tenantId, tenantId); } @@ -141,24 +145,30 @@ public class BaseQueueService extends AbstractEntityService implements QueueServ return Optional.ofNullable(findQueueById(tenantId, new QueueId(entityId.getId()))); } + @Override + public FluentFuture>> findEntityAsync(TenantId tenantId, EntityId entityId) { + return FluentFuture.from(queueDao.findByIdAsync(tenantId, entityId.getId())) + .transform(Optional::ofNullable, directExecutor()); + } + @Override public EntityType getEntityType() { return EntityType.QUEUE; } - private PaginatedRemover tenantQueuesRemover = - new PaginatedRemover<>() { + private final PaginatedRemover tenantQueuesRemover = new PaginatedRemover<>() { - @Override - protected PageData findEntities(TenantId tenantId, TenantId id, PageLink pageLink) { - return queueDao.findQueuesByTenantId(id, pageLink); - } + @Override + protected PageData findEntities(TenantId tenantId, TenantId id, PageLink pageLink) { + return queueDao.findQueuesByTenantId(id, pageLink); + } - @Override - protected void removeEntity(TenantId tenantId, Queue entity) { - deleteQueue(tenantId, entity.getId()); - } - }; + @Override + protected void removeEntity(TenantId tenantId, Queue entity) { + deleteQueue(tenantId, entity.getId()); + } + + }; private TenantId getSystemOrIsolatedTenantId(TenantId tenantId) { if (!tenantId.equals(TenantId.SYS_TENANT_ID)) { @@ -167,7 +177,7 @@ public class BaseQueueService extends AbstractEntityService implements QueueServ return tenantId; } } - return TenantId.SYS_TENANT_ID; } + } diff --git a/dao/src/main/java/org/thingsboard/server/dao/queue/BaseQueueStatsService.java b/dao/src/main/java/org/thingsboard/server/dao/queue/BaseQueueStatsService.java index 9e4c7136d5..4091daac20 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/queue/BaseQueueStatsService.java +++ b/dao/src/main/java/org/thingsboard/server/dao/queue/BaseQueueStatsService.java @@ -15,6 +15,7 @@ */ package org.thingsboard.server.dao.queue; +import com.google.common.util.concurrent.FluentFuture; import lombok.RequiredArgsConstructor; import lombok.extern.slf4j.Slf4j; import org.springframework.stereotype.Service; @@ -35,6 +36,7 @@ import org.thingsboard.server.dao.service.Validator; import java.util.List; import java.util.Optional; +import static com.google.common.util.concurrent.MoreExecutors.directExecutor; import static org.thingsboard.server.dao.service.Validator.validateId; import static org.thingsboard.server.dao.service.Validator.validateIds; @@ -106,6 +108,12 @@ public class BaseQueueStatsService extends AbstractEntityService implements Queu return Optional.ofNullable(findQueueStatsById(tenantId, new QueueStatsId(entityId.getId()))); } + @Override + public FluentFuture>> findEntityAsync(TenantId tenantId, EntityId entityId) { + return FluentFuture.from(queueStatsDao.findByIdAsync(tenantId, entityId.getId())) + .transform(Optional::ofNullable, directExecutor()); + } + @Override public EntityType getEntityType() { return EntityType.QUEUE_STATS; diff --git a/dao/src/main/java/org/thingsboard/server/dao/resource/BaseImageService.java b/dao/src/main/java/org/thingsboard/server/dao/resource/BaseImageService.java index c16e37d30f..25ef83d556 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/resource/BaseImageService.java +++ b/dao/src/main/java/org/thingsboard/server/dao/resource/BaseImageService.java @@ -17,11 +17,11 @@ package org.thingsboard.server.dao.resource; import com.fasterxml.jackson.databind.JsonNode; import jakarta.annotation.PostConstruct; -import lombok.Data; import lombok.SneakyThrows; import lombok.extern.slf4j.Slf4j; import org.apache.commons.lang3.RandomStringUtils; import org.apache.commons.lang3.StringUtils; +import org.apache.commons.lang3.Strings; import org.apache.commons.lang3.exception.ExceptionUtils; import org.apache.commons.lang3.tuple.Pair; import org.springframework.stereotype.Service; @@ -356,8 +356,8 @@ public class BaseImageService extends BaseResourceService implements ImageServic imageName = imageName + type + " image"; UpdateResult result = convertToImageUrl(entity.getTenantId(), imageName, entity.getImage(), Collections.emptyMap()); - entity.setImage(result.getValue()); - return result.isUpdated(); + entity.setImage(result.value()); + return result.updated(); } @Transactional(noRollbackFor = Exception.class) // we don't want transaction to rollback in case of an image processing failure @@ -373,8 +373,8 @@ public class BaseImageService extends BaseResourceService implements ImageServic Map imagesLinks = getResourcesLinks(widgetTypeDetails.getResources()); UpdateResult result = convertToImageUrl(tenantId, prefix + " image", widgetTypeDetails.getImage(), imagesLinks); - boolean updated = result.isUpdated(); - widgetTypeDetails.setImage(result.getValue()); + boolean updated = result.updated(); + widgetTypeDetails.setImage(result.value()); if (widgetTypeDetails.getDescriptor().isObject()) { JsonNode defaultConfig = widgetTypeDetails.getDefaultConfig(); @@ -397,8 +397,8 @@ public class BaseImageService extends BaseResourceService implements ImageServic Map imagesLinks = getResourcesLinks(dashboard.getResources()); var result = convertToImageUrl(tenantId, prefix + " image", dashboard.getImage(), imagesLinks); - boolean updated = result.isUpdated(); - dashboard.setImage(result.getValue()); + boolean updated = result.updated(); + dashboard.setImage(result.value()); updated |= convertToImageUrlsByMapping(tenantId, DASHBOARD_BASE64_MAPPING, Collections.singletonMap("prefix", prefix), dashboard.getConfiguration(), imagesLinks); updated |= convertToImageUrls(tenantId, prefix, dashboard.getConfiguration(), imagesLinks); @@ -409,10 +409,10 @@ public class BaseImageService extends BaseResourceService implements ImageServic AtomicBoolean updated = new AtomicBoolean(false); JacksonUtil.replaceAllByMapping(configuration, mapping, templateParams, (name, value) -> { UpdateResult result = convertToImageUrl(tenantId, name, value, links); - if (result.isUpdated()) { + if (result.updated()) { updated.set(true); } - return result.getValue(); + return result.value(); }); return updated.get(); } @@ -514,10 +514,10 @@ public class BaseImageService extends BaseResourceService implements ImageServic AtomicBoolean updated = new AtomicBoolean(false); JacksonUtil.replaceAll(root, title, (path, value) -> { UpdateResult result = convertToImageUrl(tenantId, path, value, true, links); - if (result.isUpdated()) { + if (result.updated()) { updated.set(true); } - return result.getValue(); + return result.value(); }); return updated.get(); } @@ -678,16 +678,18 @@ public class BaseImageService extends BaseResourceService implements ImageServic private String getImageLink(String value) { if (value.startsWith(DataConstants.TB_IMAGE_PREFIX + "/api/images")) { - return StringUtils.removeStart(value, DataConstants.TB_IMAGE_PREFIX); + return Strings.CS.removeStart(value, DataConstants.TB_IMAGE_PREFIX); } else { return null; } } - @Data(staticConstructor = "of") - private static class UpdateResult { - private final boolean updated; - private final String value; + private record UpdateResult(boolean updated, String value) { + + static UpdateResult of(boolean updated, String value) { + return new UpdateResult(updated, value); + } + } } diff --git a/dao/src/main/java/org/thingsboard/server/dao/resource/BaseResourceService.java b/dao/src/main/java/org/thingsboard/server/dao/resource/BaseResourceService.java index 7f0dd4a6bf..1eed9a6b59 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/resource/BaseResourceService.java +++ b/dao/src/main/java/org/thingsboard/server/dao/resource/BaseResourceService.java @@ -18,6 +18,7 @@ package org.thingsboard.server.dao.resource; import com.fasterxml.jackson.databind.JsonNode; import com.fasterxml.jackson.databind.node.TextNode; import com.google.common.hash.Hashing; +import com.google.common.util.concurrent.FluentFuture; import com.google.common.util.concurrent.ListenableFuture; import jakarta.annotation.PostConstruct; import lombok.RequiredArgsConstructor; @@ -79,6 +80,7 @@ import java.util.UUID; import java.util.concurrent.atomic.AtomicBoolean; import java.util.function.UnaryOperator; +import static com.google.common.util.concurrent.MoreExecutors.directExecutor; import static org.thingsboard.server.common.data.StringUtils.isNotEmpty; import static org.thingsboard.server.dao.device.DeviceServiceImpl.INCORRECT_TENANT_ID; import static org.thingsboard.server.dao.service.Validator.validateId; @@ -107,7 +109,8 @@ public class BaseResourceService extends AbstractCachedEntityService DASHBOARD_RESOURCES_MAPPING = Map.of( @@ -275,7 +278,7 @@ public class BaseResourceService extends AbstractCachedEntityService>> findEntityAsync(TenantId tenantId, EntityId entityId) { + return FluentFuture.from(findResourceInfoByIdAsync(tenantId, new TbResourceId(entityId.getId()))) + .transform(Optional::ofNullable, directExecutor()); + } + @Override public EntityType getEntityType() { return EntityType.TB_RESOURCE; diff --git a/dao/src/main/java/org/thingsboard/server/dao/rpc/BaseRpcService.java b/dao/src/main/java/org/thingsboard/server/dao/rpc/BaseRpcService.java index 15b9717f47..f06e2e45df 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/rpc/BaseRpcService.java +++ b/dao/src/main/java/org/thingsboard/server/dao/rpc/BaseRpcService.java @@ -15,6 +15,7 @@ */ package org.thingsboard.server.dao.rpc; +import com.google.common.util.concurrent.FluentFuture; import com.google.common.util.concurrent.ListenableFuture; import lombok.RequiredArgsConstructor; import lombok.extern.slf4j.Slf4j; @@ -33,6 +34,7 @@ import org.thingsboard.server.dao.service.PaginatedRemover; import java.util.Optional; +import static com.google.common.util.concurrent.MoreExecutors.directExecutor; import static org.thingsboard.server.dao.service.Validator.validateId; import static org.thingsboard.server.dao.service.Validator.validatePageLink; @@ -40,6 +42,7 @@ import static org.thingsboard.server.dao.service.Validator.validatePageLink; @Slf4j @RequiredArgsConstructor public class BaseRpcService implements RpcService { + public static final String INCORRECT_TENANT_ID = "Incorrect tenantId "; public static final String INCORRECT_RPC_ID = "Incorrect rpcId "; @@ -113,21 +116,29 @@ public class BaseRpcService implements RpcService { return Optional.ofNullable(findById(tenantId, new RpcId(entityId.getId()))); } + @Override + public FluentFuture>> findEntityAsync(TenantId tenantId, EntityId entityId) { + return FluentFuture.from(findRpcByIdAsync(tenantId, new RpcId(entityId.getId()))) + .transform(Optional::ofNullable, directExecutor()); + } + @Override public EntityType getEntityType() { return EntityType.RPC; } - private PaginatedRemover tenantRpcRemover = - new PaginatedRemover<>() { - @Override - protected PageData findEntities(TenantId tenantId, TenantId id, PageLink pageLink) { - return rpcDao.findAllRpcByTenantId(id, pageLink); - } - - @Override - protected void removeEntity(TenantId tenantId, Rpc entity) { - deleteRpc(tenantId, entity.getId()); - } - }; + private final PaginatedRemover tenantRpcRemover = new PaginatedRemover<>() { + + @Override + protected PageData findEntities(TenantId tenantId, TenantId id, PageLink pageLink) { + return rpcDao.findAllRpcByTenantId(id, pageLink); + } + + @Override + protected void removeEntity(TenantId tenantId, Rpc entity) { + deleteRpc(tenantId, entity.getId()); + } + + }; + } diff --git a/dao/src/main/java/org/thingsboard/server/dao/rule/BaseRuleChainService.java b/dao/src/main/java/org/thingsboard/server/dao/rule/BaseRuleChainService.java index 8538bd9492..8912e47c5b 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/rule/BaseRuleChainService.java +++ b/dao/src/main/java/org/thingsboard/server/dao/rule/BaseRuleChainService.java @@ -18,6 +18,7 @@ package org.thingsboard.server.dao.rule; import com.datastax.oss.driver.api.core.uuid.Uuids; import com.fasterxml.jackson.databind.JsonNode; import com.fasterxml.jackson.databind.node.ObjectNode; +import com.google.common.util.concurrent.FluentFuture; import com.google.common.util.concurrent.ListenableFuture; import lombok.extern.slf4j.Slf4j; import org.apache.commons.collections4.CollectionUtils; @@ -28,6 +29,7 @@ import org.springframework.stereotype.Service; import org.springframework.transaction.annotation.Transactional; import org.thingsboard.common.util.JacksonUtil; import org.thingsboard.server.common.data.BaseData; +import org.thingsboard.server.common.data.BaseDataWithAdditionalInfo; import org.thingsboard.server.common.data.EntityType; import org.thingsboard.server.common.data.audit.ActionType; import org.thingsboard.server.common.data.edge.Edge; @@ -38,6 +40,7 @@ import org.thingsboard.server.common.data.id.HasId; 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.id.UUIDBased; import org.thingsboard.server.common.data.page.PageData; import org.thingsboard.server.common.data.page.PageDataIterable; import org.thingsboard.server.common.data.page.PageLink; @@ -80,6 +83,7 @@ import java.util.UUID; import java.util.function.Function; import java.util.stream.Collectors; +import static com.google.common.util.concurrent.MoreExecutors.directExecutor; import static org.thingsboard.server.common.data.DataConstants.TENANT; import static org.thingsboard.server.dao.service.Validator.validateId; import static org.thingsboard.server.dao.service.Validator.validateIds; @@ -87,9 +91,6 @@ import static org.thingsboard.server.dao.service.Validator.validatePageLink; import static org.thingsboard.server.dao.service.Validator.validatePositiveNumber; import static org.thingsboard.server.dao.service.Validator.validateString; -/** - * Created by igor on 3/12/18. - */ @Service("RuleChainDaoService") @Slf4j public class BaseRuleChainService extends AbstractEntityService implements RuleChainService { @@ -98,6 +99,7 @@ public class BaseRuleChainService extends AbstractEntityService implements RuleC public static final String INCORRECT_TENANT_ID = "Incorrect tenantId "; public static final String TB_RULE_CHAIN_INPUT_NODE = "org.thingsboard.rule.engine.flow.TbRuleChainInputNode"; + @Autowired private RuleChainDao ruleChainDao; @@ -260,7 +262,7 @@ public class BaseRuleChainService extends AbstractEntityService implements RuleC firstRuleNodeId = nodes.get(ruleChainMetaData.getFirstNodeIndex()).getId(); } if ((ruleChain.getFirstRuleNodeId() != null && !ruleChain.getFirstRuleNodeId().equals(firstRuleNodeId)) - || (ruleChain.getFirstRuleNodeId() == null && firstRuleNodeId != null)) { + || (ruleChain.getFirstRuleNodeId() == null && firstRuleNodeId != null)) { ruleChain.setFirstRuleNodeId(firstRuleNodeId); } if (ruleChainMetaData.getConnections() != null) { @@ -876,6 +878,17 @@ public class BaseRuleChainService extends AbstractEntityService implements RuleC return Optional.ofNullable(hasId); } + @Override + public FluentFuture>> findEntityAsync(TenantId tenantId, EntityId entityId) { + ListenableFuture> future; + if (entityId.getEntityType() == EntityType.RULE_NODE) { + future = findRuleNodeByIdAsync(tenantId, new RuleNodeId(entityId.getId())); + } else { + future = findRuleChainByIdAsync(tenantId, new RuleChainId(entityId.getId())); + } + return FluentFuture.from(future).transform(Optional::ofNullable, directExecutor()); + } + @Override public long countByTenantId(TenantId tenantId) { return ruleChainDao.countByTenantId(tenantId); @@ -919,18 +932,11 @@ public class BaseRuleChainService extends AbstractEntityService implements RuleC ComponentClusteringMode nodeConfigType = ReflectionUtils.getAnnotationProperty(ruleNode.getType(), "org.thingsboard.rule.engine.api.RuleNode", "clusteringMode"); - switch (nodeConfigType) { - case ENABLED: - singletonMode = false; - break; - case SINGLETON: - singletonMode = true; - break; - case USER_PREFERENCE: - default: - singletonMode = ruleNode.isSingletonMode(); - break; - } + singletonMode = switch (nodeConfigType) { + case ENABLED -> false; + case SINGLETON -> true; + default -> ruleNode.isSingletonMode(); + }; } catch (Exception e) { log.warn("Failed to get clustering mode: {}", ExceptionUtils.getRootCauseMessage(e)); singletonMode = false; diff --git a/dao/src/main/java/org/thingsboard/server/dao/settings/AdminSettingsServiceImpl.java b/dao/src/main/java/org/thingsboard/server/dao/settings/AdminSettingsServiceImpl.java index 604a37f4bb..7440481700 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/settings/AdminSettingsServiceImpl.java +++ b/dao/src/main/java/org/thingsboard/server/dao/settings/AdminSettingsServiceImpl.java @@ -17,6 +17,7 @@ package org.thingsboard.server.dao.settings; import com.fasterxml.jackson.databind.JsonNode; import com.fasterxml.jackson.databind.node.ObjectNode; +import com.google.common.util.concurrent.FluentFuture; import lombok.extern.slf4j.Slf4j; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.stereotype.Service; @@ -31,6 +32,8 @@ import org.thingsboard.server.dao.service.Validator; import java.util.Optional; +import static com.google.common.util.concurrent.MoreExecutors.directExecutor; + @Service @Slf4j public class AdminSettingsServiceImpl implements AdminSettingsService { @@ -106,6 +109,12 @@ public class AdminSettingsServiceImpl implements AdminSettingsService { return Optional.ofNullable(adminSettingsDao.findById(tenantId, entityId.getId())); } + @Override + public FluentFuture>> findEntityAsync(TenantId tenantId, EntityId entityId) { + return FluentFuture.from(adminSettingsDao.findByIdAsync(tenantId, entityId.getId())) + .transform(Optional::ofNullable, directExecutor()); + } + @Override public EntityType getEntityType() { return EntityType.ADMIN_SETTINGS; diff --git a/dao/src/main/java/org/thingsboard/server/dao/tenant/TenantProfileServiceImpl.java b/dao/src/main/java/org/thingsboard/server/dao/tenant/TenantProfileServiceImpl.java index e165d6cfe6..b9172306bb 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/tenant/TenantProfileServiceImpl.java +++ b/dao/src/main/java/org/thingsboard/server/dao/tenant/TenantProfileServiceImpl.java @@ -15,6 +15,7 @@ */ package org.thingsboard.server.dao.tenant; +import com.google.common.util.concurrent.FluentFuture; import lombok.extern.slf4j.Slf4j; import org.hibernate.exception.ConstraintViolationException; import org.springframework.beans.factory.annotation.Autowired; @@ -44,6 +45,7 @@ import java.util.List; import java.util.Optional; import java.util.UUID; +import static com.google.common.util.concurrent.MoreExecutors.directExecutor; import static org.thingsboard.common.util.DebugModeUtil.DEBUG_MODE_DEFAULT_DURATION_MINUTES; import static org.thingsboard.server.dao.service.Validator.validateId; @@ -230,23 +232,29 @@ public class TenantProfileServiceImpl extends AbstractCachedEntityService>> findEntityAsync(TenantId tenantId, EntityId entityId) { + return FluentFuture.from(tenantProfileDao.findByIdAsync(tenantId, entityId.getId())) + .transform(Optional::ofNullable, directExecutor()); + } + @Override public EntityType getEntityType() { return EntityType.TENANT_PROFILE; } - private final PaginatedRemover tenantProfilesRemover = - new PaginatedRemover<>() { + private final PaginatedRemover tenantProfilesRemover = new PaginatedRemover<>() { - @Override - protected PageData findEntities(TenantId tenantId, String id, PageLink pageLink) { - return tenantProfileDao.findTenantProfiles(tenantId, pageLink); - } + @Override + protected PageData findEntities(TenantId tenantId, String id, PageLink pageLink) { + return tenantProfileDao.findTenantProfiles(tenantId, pageLink); + } + + @Override + protected void removeEntity(TenantId tenantId, TenantProfile entity) { + removeTenantProfile(tenantId, entity, entity.isDefault()); + } - @Override - protected void removeEntity(TenantId tenantId, TenantProfile entity) { - removeTenantProfile(tenantId, entity, entity.isDefault()); - } - }; + }; } diff --git a/dao/src/main/java/org/thingsboard/server/dao/tenant/TenantServiceImpl.java b/dao/src/main/java/org/thingsboard/server/dao/tenant/TenantServiceImpl.java index 0df7c36527..bd5a621562 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/tenant/TenantServiceImpl.java +++ b/dao/src/main/java/org/thingsboard/server/dao/tenant/TenantServiceImpl.java @@ -15,6 +15,7 @@ */ package org.thingsboard.server.dao.tenant; +import com.google.common.util.concurrent.FluentFuture; import com.google.common.util.concurrent.ListenableFuture; import lombok.extern.slf4j.Slf4j; import org.springframework.beans.factory.annotation.Autowired; @@ -51,6 +52,7 @@ import java.util.List; import java.util.Optional; import java.util.function.Consumer; +import static com.google.common.util.concurrent.MoreExecutors.directExecutor; import static org.thingsboard.server.dao.service.Validator.validateId; @Service("TenantDaoService") @@ -85,8 +87,8 @@ public class TenantServiceImpl extends AbstractCachedEntityService existsTenantCache; - @TransactionalEventListener(classes = TenantEvictEvent.class) @Override + @TransactionalEventListener public void handleEvictEvent(TenantEvictEvent event) { TenantId tenantId = event.getTenantId(); cache.evict(tenantId); @@ -245,6 +247,12 @@ public class TenantServiceImpl extends AbstractCachedEntityService>> findEntityAsync(TenantId tenantId, EntityId entityId) { + return FluentFuture.from(findTenantByIdAsync(tenantId, TenantId.fromUUID(entityId.getId()))) + .transform(Optional::ofNullable, directExecutor()); + } + @Override public EntityType getEntityType() { return EntityType.TENANT; diff --git a/dao/src/main/java/org/thingsboard/server/dao/usagerecord/ApiUsageStateServiceImpl.java b/dao/src/main/java/org/thingsboard/server/dao/usagerecord/ApiUsageStateServiceImpl.java index 1282493650..49ee56e7be 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/usagerecord/ApiUsageStateServiceImpl.java +++ b/dao/src/main/java/org/thingsboard/server/dao/usagerecord/ApiUsageStateServiceImpl.java @@ -15,6 +15,7 @@ */ package org.thingsboard.server.dao.usagerecord; +import com.google.common.util.concurrent.FluentFuture; import lombok.extern.slf4j.Slf4j; import org.springframework.context.annotation.Lazy; import org.springframework.stereotype.Service; @@ -48,11 +49,13 @@ import java.util.List; import java.util.Objects; import java.util.Optional; +import static com.google.common.util.concurrent.MoreExecutors.directExecutor; import static org.thingsboard.server.dao.service.Validator.validateId; @Service("ApiUsageStateDaoService") @Slf4j public class ApiUsageStateServiceImpl extends AbstractEntityService implements ApiUsageStateService { + public static final String INCORRECT_TENANT_ID = "Incorrect tenantId "; private final ApiUsageStateDao apiUsageStateDao; @@ -161,7 +164,7 @@ public class ApiUsageStateServiceImpl extends AbstractEntityService implements A public ApiUsageState update(ApiUsageState apiUsageState) { log.trace("Executing save [{}]", apiUsageState.getTenantId()); validateId(apiUsageState.getTenantId(), id -> INCORRECT_TENANT_ID + id); - validateId(apiUsageState.getId(), "Can't save new usage state. Only update is allowed!"); + validateId(apiUsageState.getId(), __ -> "Can't save new usage state. Only update is allowed!"); apiUsageState.setVersion(null); ApiUsageState savedState = apiUsageStateDao.save(apiUsageState.getTenantId(), apiUsageState); eventPublisher.publishEvent(SaveEntityEvent.builder().tenantId(savedState.getTenantId()).entityId(savedState.getId()) @@ -195,6 +198,12 @@ public class ApiUsageStateServiceImpl extends AbstractEntityService implements A return Optional.ofNullable(findApiUsageStateById(tenantId, new ApiUsageStateId(entityId.getId()))); } + @Override + public FluentFuture>> findEntityAsync(TenantId tenantId, EntityId entityId) { + return FluentFuture.from(apiUsageStateDao.findByIdAsync(tenantId, entityId.getId())) + .transform(Optional::ofNullable, directExecutor()); + } + @Transactional @Override public void deleteByTenantId(TenantId tenantId) { diff --git a/dao/src/main/java/org/thingsboard/server/dao/user/UserServiceImpl.java b/dao/src/main/java/org/thingsboard/server/dao/user/UserServiceImpl.java index 71a908189f..977bb01785 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/user/UserServiceImpl.java +++ b/dao/src/main/java/org/thingsboard/server/dao/user/UserServiceImpl.java @@ -18,6 +18,7 @@ package org.thingsboard.server.dao.user; import com.fasterxml.jackson.core.type.TypeReference; import com.fasterxml.jackson.databind.JsonNode; import com.fasterxml.jackson.databind.node.ObjectNode; +import com.google.common.util.concurrent.FluentFuture; import com.google.common.util.concurrent.ListenableFuture; import lombok.RequiredArgsConstructor; import lombok.extern.slf4j.Slf4j; @@ -72,6 +73,7 @@ import java.util.Objects; import java.util.Optional; import java.util.concurrent.TimeUnit; +import static com.google.common.util.concurrent.MoreExecutors.directExecutor; import static org.thingsboard.server.common.data.StringUtils.generateSafeToken; import static org.thingsboard.server.dao.service.Validator.validateId; import static org.thingsboard.server.dao.service.Validator.validatePageLink; @@ -103,8 +105,8 @@ public class UserServiceImpl extends AbstractCachedEntityService keys = new ArrayList<>(2); keys.add(new UserCacheKey(event.tenantId(), event.newEmail())); @@ -568,6 +570,12 @@ public class UserServiceImpl extends AbstractCachedEntityService>> findEntityAsync(TenantId tenantId, EntityId entityId) { + return FluentFuture.from(findUserByIdAsync(tenantId, new UserId(entityId.getId()))) + .transform(Optional::ofNullable, directExecutor()); + } + @Override public long countByTenantId(TenantId tenantId) { return userDao.countByTenantId(tenantId); diff --git a/dao/src/main/java/org/thingsboard/server/dao/widget/WidgetTypeServiceImpl.java b/dao/src/main/java/org/thingsboard/server/dao/widget/WidgetTypeServiceImpl.java index 7e74c2a18b..db8cd875e3 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/widget/WidgetTypeServiceImpl.java +++ b/dao/src/main/java/org/thingsboard/server/dao/widget/WidgetTypeServiceImpl.java @@ -15,6 +15,7 @@ */ package org.thingsboard.server.dao.widget; +import com.google.common.util.concurrent.FluentFuture; import lombok.extern.slf4j.Slf4j; import org.apache.commons.collections4.CollectionUtils; import org.springframework.beans.factory.annotation.Autowired; @@ -47,6 +48,7 @@ import java.util.ArrayList; import java.util.List; import java.util.Optional; +import static com.google.common.util.concurrent.MoreExecutors.directExecutor; import static org.thingsboard.server.dao.service.Validator.validateIds; @Service("WidgetTypeDaoService") @@ -54,8 +56,6 @@ import static org.thingsboard.server.dao.service.Validator.validateIds; public class WidgetTypeServiceImpl implements WidgetTypeService { public static final String INCORRECT_TENANT_ID = "Incorrect tenantId "; - public static final String INCORRECT_RESOURCE_ID = "Incorrect resourceId "; - public static final String INCORRECT_BUNDLE_ALIAS = "Incorrect bundleAlias "; public static final String INCORRECT_WIDGETS_BUNDLE_ID = "Incorrect widgetsBundleId "; @Autowired @@ -281,6 +281,12 @@ public class WidgetTypeServiceImpl implements WidgetTypeService { return Optional.ofNullable(findWidgetTypeById(tenantId, new WidgetTypeId(entityId.getId()))); } + @Override + public FluentFuture>> findEntityAsync(TenantId tenantId, EntityId entityId) { + return FluentFuture.from(widgetTypeDao.findByIdAsync(tenantId, entityId.getId())) + .transform(Optional::ofNullable, directExecutor()); + } + @Override public EntityType getEntityType() { return EntityType.WIDGET_TYPE; @@ -303,6 +309,7 @@ public class WidgetTypeServiceImpl implements WidgetTypeService { protected void removeEntity(TenantId tenantId, WidgetTypeInfo entity) { deleteWidgetType(tenantId, new WidgetTypeId(entity.getUuidId())); } + }; private final PaginatedRemover bundleWidgetTypesRemover = new PaginatedRemover<>() { @@ -316,6 +323,7 @@ public class WidgetTypeServiceImpl implements WidgetTypeService { protected void removeEntity(TenantId tenantId, WidgetTypeInfo widgetTypeInfo) { deleteWidgetType(tenantId, widgetTypeInfo.getId()); } + }; } diff --git a/dao/src/main/java/org/thingsboard/server/dao/widget/WidgetsBundleServiceImpl.java b/dao/src/main/java/org/thingsboard/server/dao/widget/WidgetsBundleServiceImpl.java index f057ce22ca..5b5930fb0a 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/widget/WidgetsBundleServiceImpl.java +++ b/dao/src/main/java/org/thingsboard/server/dao/widget/WidgetsBundleServiceImpl.java @@ -16,6 +16,7 @@ package org.thingsboard.server.dao.widget; import com.fasterxml.jackson.databind.JsonNode; +import com.google.common.util.concurrent.FluentFuture; import lombok.extern.slf4j.Slf4j; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.context.ApplicationEventPublisher; @@ -33,12 +34,10 @@ import org.thingsboard.server.common.data.widget.WidgetType; import org.thingsboard.server.common.data.widget.WidgetTypeDetails; import org.thingsboard.server.common.data.widget.WidgetsBundle; import org.thingsboard.server.common.data.widget.WidgetsBundleFilter; -import org.thingsboard.server.dao.entity.AbstractCachedEntityService; import org.thingsboard.server.dao.eventsourcing.DeleteEntityEvent; import org.thingsboard.server.dao.eventsourcing.SaveEntityEvent; import org.thingsboard.server.dao.exception.IncorrectParameterException; import org.thingsboard.server.dao.resource.ImageService; -import org.thingsboard.server.dao.resource.ResourceService; import org.thingsboard.server.dao.service.DataValidator; import org.thingsboard.server.dao.service.PaginatedRemover; import org.thingsboard.server.dao.service.Validator; @@ -48,13 +47,15 @@ import java.util.List; import java.util.Optional; import java.util.stream.Stream; +import static com.google.common.util.concurrent.MoreExecutors.directExecutor; +import static org.thingsboard.server.dao.entity.AbstractEntityService.checkConstraintViolation; + @Service("WidgetsBundleDaoService") @Slf4j public class WidgetsBundleServiceImpl implements WidgetsBundleService { private static final int DEFAULT_WIDGETS_BUNDLE_LIMIT = 300; public static final String INCORRECT_TENANT_ID = "Incorrect tenantId "; - public static final String INCORRECT_PAGE_LINK = "Incorrect page link "; @Autowired private WidgetsBundleDao widgetsBundleDao; @@ -71,9 +72,6 @@ public class WidgetsBundleServiceImpl implements WidgetsBundleService { @Autowired private ImageService imageService; - @Autowired - private ResourceService resourceService; - @Override public WidgetsBundle findWidgetsBundleById(TenantId tenantId, WidgetsBundleId widgetsBundleId) { log.trace("Executing findWidgetsBundleById [{}]", widgetsBundleId); @@ -92,9 +90,9 @@ public class WidgetsBundleServiceImpl implements WidgetsBundleService { .entityId(result.getId()).created(widgetsBundle.getId() == null).build()); return result; } catch (Exception e) { - AbstractCachedEntityService.checkConstraintViolation(e, - "uq_widgets_bundle_alias", "Widgets Bundle with such alias already exists!"); - AbstractCachedEntityService.checkConstraintViolation(e, "widgets_bundle_external_id_unq_key", "Widgets Bundle with such external id already exists!"); + checkConstraintViolation(e, + "uq_widgets_bundle_alias", "Widgets Bundle with such alias already exists!", + "widgets_bundle_external_id_unq_key", "Widgets Bundle with such external id already exists!"); throw e; } } @@ -248,9 +246,7 @@ public class WidgetsBundleServiceImpl implements WidgetsBundleService { } if (widgetsBundleDescriptor.has("widgetTypeFqns")) { JsonNode widgetFqnsArrayJson = widgetsBundleDescriptor.get("widgetTypeFqns"); - widgetFqnsArrayJson.forEach(fqnJson -> { - widgetTypeFqns.add(fqnJson.asText()); - }); + widgetFqnsArrayJson.forEach(fqnJson -> widgetTypeFqns.add(fqnJson.asText())); } widgetTypeService.updateWidgetsBundleWidgetFqns(TenantId.SYS_TENANT_ID, widgetsBundle.getId(), widgetTypeFqns); }); @@ -274,23 +270,29 @@ public class WidgetsBundleServiceImpl implements WidgetsBundleService { return Optional.ofNullable(findWidgetsBundleById(tenantId, new WidgetsBundleId(entityId.getId()))); } + @Override + public FluentFuture>> findEntityAsync(TenantId tenantId, EntityId entityId) { + return FluentFuture.from(widgetsBundleDao.findByIdAsync(tenantId, entityId.getId())) + .transform(Optional::ofNullable, directExecutor()); + } + @Override public EntityType getEntityType() { return EntityType.WIDGETS_BUNDLE; } - private PaginatedRemover tenantWidgetsBundleRemover = - new PaginatedRemover() { + private final PaginatedRemover tenantWidgetsBundleRemover = new PaginatedRemover<>() { - @Override - protected PageData findEntities(TenantId tenantId, TenantId id, PageLink pageLink) { - return widgetsBundleDao.findTenantWidgetsBundlesByTenantId(id.getId(), pageLink); - } + @Override + protected PageData findEntities(TenantId tenantId, TenantId id, PageLink pageLink) { + return widgetsBundleDao.findTenantWidgetsBundlesByTenantId(id.getId(), pageLink); + } + + @Override + protected void removeEntity(TenantId tenantId, WidgetsBundle entity) { + deleteWidgetsBundle(tenantId, new WidgetsBundleId(entity.getUuidId())); + } - @Override - protected void removeEntity(TenantId tenantId, WidgetsBundle entity) { - deleteWidgetsBundle(tenantId, new WidgetsBundleId(entity.getUuidId())); - } - }; + }; } diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbAbstractGetEntityDataNode.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbAbstractGetEntityDataNode.java index ebf32bf6b1..1982596e91 100644 --- a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbAbstractGetEntityDataNode.java +++ b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbAbstractGetEntityDataNode.java @@ -58,15 +58,9 @@ public abstract class TbAbstractGetEntityDataNode extends Tb protected void processDataAndTell(TbContext ctx, TbMsg msg, T entityId, ObjectNode msgDataAsJsonNode) { DataToFetch dataToFetch = config.getDataToFetch(); switch (dataToFetch) { - case ATTRIBUTES: - processAttributesKvEntryData(ctx, msg, entityId, msgDataAsJsonNode); - break; - case LATEST_TELEMETRY: - processTsKvEntryData(ctx, msg, entityId, msgDataAsJsonNode); - break; - case FIELDS: - processFieldsData(ctx, msg, entityId, msgDataAsJsonNode, true); - break; + case ATTRIBUTES -> processAttributesKvEntryData(ctx, msg, entityId, msgDataAsJsonNode); + case LATEST_TELEMETRY -> processTsKvEntryData(ctx, msg, entityId, msgDataAsJsonNode); + case FIELDS -> processFieldsData(ctx, msg, entityId, msgDataAsJsonNode, true); } } diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbAbstractNodeWithFetchTo.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbAbstractNodeWithFetchTo.java index 189ec47ef1..a11e3333ff 100644 --- a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbAbstractNodeWithFetchTo.java +++ b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbAbstractNodeWithFetchTo.java @@ -17,8 +17,7 @@ package org.thingsboard.rule.engine.metadata; import com.fasterxml.jackson.databind.JsonNode; import com.fasterxml.jackson.databind.node.ObjectNode; -import com.google.common.util.concurrent.AsyncFunction; -import com.google.common.util.concurrent.Futures; +import com.google.common.base.Function; import lombok.extern.slf4j.Slf4j; import org.thingsboard.common.util.JacksonUtil; import org.thingsboard.rule.engine.api.TbContext; @@ -54,12 +53,12 @@ public abstract class TbAbstractNodeWithFetchTo AsyncFunction checkIfEntityIsPresentOrThrow(String message) { + protected Function checkIfEntityIsPresentOrThrow(String message) { return id -> { if (id == null || id.isNullUid()) { - return Futures.immediateFailedFuture(new NoSuchElementException(message)); + throw new NoSuchElementException(message); } - return Futures.immediateFuture(id); + return id; }; } diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetCustomerAttributeNode.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetCustomerAttributeNode.java index d779bacb77..02a1330a53 100644 --- a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetCustomerAttributeNode.java +++ b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetCustomerAttributeNode.java @@ -16,19 +16,22 @@ 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.ListenableFuture; import org.thingsboard.rule.engine.api.RuleNode; 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.api.util.TbNodeUtils; -import org.thingsboard.rule.engine.util.EntitiesCustomerIdAsyncLoader; +import org.thingsboard.server.common.data.EntityType; import org.thingsboard.server.common.data.id.CustomerId; import org.thingsboard.server.common.data.id.EntityId; import org.thingsboard.server.common.data.plugin.ComponentType; import org.thingsboard.server.common.data.util.TbPair; +import java.util.NoSuchElementException; + +import static com.google.common.util.concurrent.Futures.immediateFuture; + @RuleNode( type = ComponentType.ENRICHMENT, name = "customer attributes", @@ -44,8 +47,6 @@ import org.thingsboard.server.common.data.util.TbPair; ) public class TbGetCustomerAttributeNode extends TbAbstractGetEntityDataNode { - private static final String CUSTOMER_NOT_FOUND_MESSAGE = "Failed to find customer for entity with id: %s and type: %s"; - @Override protected TbGetEntityDataNodeConfiguration loadNodeConfiguration(TbNodeConfiguration configuration) throws TbNodeException { var config = TbNodeUtils.convert(configuration, TbGetEntityDataNodeConfiguration.class); @@ -56,10 +57,19 @@ public class TbGetCustomerAttributeNode extends TbAbstractGetEntityDataNode findEntityAsync(TbContext ctx, EntityId originator) { - return Futures.transformAsync(EntitiesCustomerIdAsyncLoader.findEntityIdAsync(ctx, originator), - checkIfEntityIsPresentOrThrow(String.format(CUSTOMER_NOT_FOUND_MESSAGE, originator.getId(), originator.getEntityType().getNormalName())), - ctx.getDbCallbackExecutor() - ); + if (originator.getEntityType() == EntityType.CUSTOMER) { + return immediateFuture((CustomerId) originator); + } + return ctx.getEntityService().fetchEntityCustomerIdAsync(ctx.getTenantId(), originator) + .transform(customerIdOpt -> { + if (customerIdOpt.isEmpty()) { + throw new NoSuchElementException("Originator not found"); + } + if (customerIdOpt.get().isNullUid()) { + throw new IllegalStateException("Originator is not assigned to any customer"); + } + return customerIdOpt.get(); + }, ctx.getDbCallbackExecutor()); } @Override diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetDeviceAttrNode.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetDeviceAttrNode.java index cf370114c3..9dbdba8ebe 100644 --- a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetDeviceAttrNode.java +++ b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetDeviceAttrNode.java @@ -55,7 +55,7 @@ public class TbGetDeviceAttrNode extends TbAbstractGetAttributesNode findEntityIdAsync(TbContext ctx, TbMsg msg) { - return Futures.transformAsync( + return Futures.transform( EntitiesRelatedDeviceIdAsyncLoader.findDeviceAsync(ctx, msg.getOriginator(), config.getDeviceRelationsQuery()), checkIfEntityIsPresentOrThrow(RELATED_DEVICE_NOT_FOUND_MESSAGE), ctx.getDbCallbackExecutor()); diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetRelatedAttributeNode.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetRelatedAttributeNode.java index dffcaf35ba..fcd11a08f1 100644 --- a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetRelatedAttributeNode.java +++ b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetRelatedAttributeNode.java @@ -58,7 +58,7 @@ public class TbGetRelatedAttributeNode extends TbAbstractGetEntityDataNode findEntityAsync(TbContext ctx, EntityId originator) { var relatedAttrConfig = (TbGetRelatedDataNodeConfiguration) config; - return Futures.transformAsync( + return Futures.transform( EntitiesRelatedEntityIdAsyncLoader.findEntityAsync(ctx, originator, relatedAttrConfig.getRelationsQuery()), checkIfEntityIsPresentOrThrow(RELATED_ENTITY_NOT_FOUND_MESSAGE), ctx.getDbCallbackExecutor()); diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/transform/TbChangeOriginatorNode.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/transform/TbChangeOriginatorNode.java index e3ab60a6e0..a97cc22d5b 100644 --- a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/transform/TbChangeOriginatorNode.java +++ b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/transform/TbChangeOriginatorNode.java @@ -25,7 +25,6 @@ import org.thingsboard.rule.engine.api.TbNodeException; import org.thingsboard.rule.engine.api.util.TbNodeUtils; import org.thingsboard.rule.engine.util.EntitiesAlarmOriginatorIdAsyncLoader; import org.thingsboard.rule.engine.util.EntitiesByNameAndTypeLoader; -import org.thingsboard.rule.engine.util.EntitiesCustomerIdAsyncLoader; import org.thingsboard.rule.engine.util.EntitiesRelatedEntityIdAsyncLoader; import org.thingsboard.server.common.data.EntityType; import org.thingsboard.server.common.data.StringUtils; @@ -36,6 +35,7 @@ import org.thingsboard.server.common.msg.TbMsg; import java.util.List; import java.util.NoSuchElementException; +import static com.google.common.util.concurrent.Futures.immediateFuture; import static org.thingsboard.rule.engine.transform.OriginatorSource.ENTITY; import static org.thingsboard.rule.engine.transform.OriginatorSource.RELATED; @@ -74,16 +74,20 @@ public class TbChangeOriginatorNode extends TbAbstractTransformNode getNewOriginator(TbContext ctx, TbMsg msg) { switch (config.getOriginatorSource()) { case CUSTOMER: - return EntitiesCustomerIdAsyncLoader.findEntityIdAsync(ctx, msg.getOriginator()); + if (msg.getOriginator().getEntityType() == EntityType.CUSTOMER) { + return immediateFuture(msg.getOriginator()); + } + return ctx.getEntityService().fetchEntityCustomerIdAsync(ctx.getTenantId(), msg.getOriginator()) + .transform(customerIdOpt -> customerIdOpt.orElse(null), ctx.getDbCallbackExecutor()); case TENANT: - return Futures.immediateFuture(ctx.getTenantId()); + return immediateFuture(ctx.getTenantId()); case RELATED: return EntitiesRelatedEntityIdAsyncLoader.findEntityAsync(ctx, msg.getOriginator(), config.getRelationsQuery()); case ALARM_ORIGINATOR: @@ -93,7 +97,7 @@ public class TbChangeOriginatorNode extends TbAbstractTransformNode findEntityIdAsync(TbContext ctx, EntityId originator) { - switch (originator.getEntityType()) { - case CUSTOMER: - return Futures.immediateFuture((CustomerId) originator); - case USER: - return toCustomerIdAsync(ctx, ctx.getUserService().findUserByIdAsync(ctx.getTenantId(), (UserId) originator)); - case ASSET: - return toCustomerIdAsync(ctx, ctx.getAssetService().findAssetByIdAsync(ctx.getTenantId(), (AssetId) originator)); - case DEVICE: - return toCustomerIdAsync(ctx, Futures.immediateFuture(ctx.getDeviceService().findDeviceById(ctx.getTenantId(), (DeviceId) originator))); - default: - return Futures.immediateFailedFuture(new TbNodeException("Unexpected originator EntityType: " + originator.getEntityType())); - } - } - - private static ListenableFuture toCustomerIdAsync(TbContext ctx, ListenableFuture future) { - return Futures.transform(future, in -> in != null ? in.getCustomerId() : null, ctx.getDbCallbackExecutor()); - } - -} diff --git a/rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/TestDbCallbackExecutor.java b/rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/TestDbCallbackExecutor.java index b3afec500c..f27ea73ce6 100644 --- a/rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/TestDbCallbackExecutor.java +++ b/rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/TestDbCallbackExecutor.java @@ -28,7 +28,7 @@ public class TestDbCallbackExecutor implements ListeningExecutor { try { return Futures.immediateFuture(task.call()); } catch (Exception e) { - throw new RuntimeException(e); + return Futures.immediateFailedFuture(e); } } diff --git a/rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/metadata/TbGetCustomerAttributeNodeTest.java b/rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/metadata/TbGetCustomerAttributeNodeTest.java index 3f2c58f191..942b1ab5ce 100644 --- a/rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/metadata/TbGetCustomerAttributeNodeTest.java +++ b/rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/metadata/TbGetCustomerAttributeNodeTest.java @@ -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( diff --git a/rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/transform/TbChangeOriginatorNodeTest.java b/rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/transform/TbChangeOriginatorNodeTest.java index 4826cfe297..da8816e824 100644 --- a/rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/transform/TbChangeOriginatorNodeTest.java +++ b/rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/transform/TbChangeOriginatorNodeTest.java @@ -15,8 +15,10 @@ */ package org.thingsboard.rule.engine.transform; +import com.google.common.util.concurrent.FluentFuture; import com.google.common.util.concurrent.Futures; import org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.api.DisplayName; import org.junit.jupiter.api.Test; import org.junit.jupiter.api.extension.ExtendWith; import org.junit.jupiter.params.ParameterizedTest; @@ -35,6 +37,7 @@ import org.thingsboard.rule.engine.api.TbNodeConfiguration; import org.thingsboard.rule.engine.api.TbNodeException; import org.thingsboard.rule.engine.api.util.TbNodeUtils; import org.thingsboard.rule.engine.data.RelationsQuery; +import org.thingsboard.server.common.data.Customer; import org.thingsboard.server.common.data.Device; import org.thingsboard.server.common.data.EntityType; import org.thingsboard.server.common.data.alarm.Alarm; @@ -54,21 +57,24 @@ import org.thingsboard.server.common.data.relation.RelationsSearchParameters; 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.device.DeviceService; +import org.thingsboard.server.dao.entity.EntityService; import org.thingsboard.server.dao.relation.RelationService; import java.util.Collections; import java.util.Map; import java.util.NoSuchElementException; +import java.util.Optional; import java.util.UUID; import java.util.stream.Stream; +import static com.google.common.util.concurrent.Futures.immediateFuture; import static org.assertj.core.api.Assertions.assertThat; import static org.assertj.core.api.Assertions.assertThatThrownBy; import static org.mockito.ArgumentMatchers.any; import static org.mockito.ArgumentMatchers.eq; import static org.mockito.BDDMockito.given; import static org.mockito.BDDMockito.then; +import static org.mockito.Mockito.lenient; import static org.thingsboard.rule.engine.transform.OriginatorSource.ALARM_ORIGINATOR; import static org.thingsboard.rule.engine.transform.OriginatorSource.CUSTOMER; import static org.thingsboard.rule.engine.transform.OriginatorSource.ENTITY; @@ -93,16 +99,23 @@ public class TbChangeOriginatorNodeTest { @Mock private AssetService assetServiceMock; @Mock - private DeviceService deviceServiceMock; - @Mock private RelationService relationServiceMock; @Mock private RuleEngineAlarmService alarmServiceMock; + @Mock + private EntityService entityServiceMock; @BeforeEach - public void before() throws TbNodeException { + public void setup() { node = new TbChangeOriginatorNode(); config = new TbChangeOriginatorNodeConfiguration().defaultConfiguration(); + + lenient().when(ctxMock.getTenantId()).thenReturn(TENANT_ID); + lenient().when(ctxMock.getDbCallbackExecutor()).thenReturn(dbExecutor); + lenient().when(ctxMock.getAssetService()).thenReturn(assetServiceMock); + lenient().when(ctxMock.getRelationService()).thenReturn(relationServiceMock); + lenient().when(ctxMock.getAlarmService()).thenReturn(alarmServiceMock); + lenient().when(ctxMock.getEntityService()).thenReturn(entityServiceMock); } @Test @@ -172,8 +185,13 @@ public class TbChangeOriginatorNodeTest { } @Test - public void givenOriginatorSourceIsCustomer_whenOnMsg_thenTellSuccess() throws TbNodeException { - Device device = new Device(DEVICE_ID); + @DisplayName(""" + Given a device assigned to a customer and node configured to change originator to customer, + when processing the message, + then should change message originator from device to customer""") + public void givenDeviceAssignedToCustomer_whenProcessingMessage_thenChangesOriginatorToCustomer() throws TbNodeException { + // GIVEN + var device = new Device(DEVICE_ID); device.setCustomerId(CUSTOMER_ID); TbMsg msg = TbMsg.newMsg() @@ -186,16 +204,52 @@ public class TbChangeOriginatorNodeTest { .originator(CUSTOMER_ID) .build(); - given(ctxMock.getDbCallbackExecutor()).willReturn(dbExecutor); - given(ctxMock.getDeviceService()).willReturn(deviceServiceMock); - given(ctxMock.getTenantId()).willReturn(TENANT_ID); - given(deviceServiceMock.findDeviceById(any(TenantId.class), any(DeviceId.class))).willReturn(device); + given(entityServiceMock.fetchEntityCustomerIdAsync(TENANT_ID, device.getId())).willReturn( + FluentFuture.from(immediateFuture(Optional.of(CUSTOMER_ID))) + ); + + given(ctxMock.transformMsgOriginator(any(TbMsg.class), any(EntityId.class))).willReturn(expectedMsg); + + node.init(ctxMock, new TbNodeConfiguration(JacksonUtil.valueToTree(config))); + + // WHEN + node.onMsg(ctxMock, msg); + + // THEN + then(ctxMock).should().transformMsgOriginator(msg, CUSTOMER_ID); + ArgumentCaptor actualMsg = ArgumentCaptor.forClass(TbMsg.class); + then(ctxMock).should().tellSuccess(actualMsg.capture()); + assertThat(actualMsg.getValue()).usingRecursiveComparison().ignoringFields("ctx").isEqualTo(expectedMsg); + } + + @Test + @DisplayName(""" + Given a customer as message originator and node configured to change originator to customer, + when processing the message, + then should keep the customer as originator""") + public void givenCustomerAsOriginator_whenProcessingMessage_thenKeepsCustomerAsOriginator() throws TbNodeException { + // GIVEN + var customer = new Customer(CUSTOMER_ID); + + var msg = TbMsg.newMsg() + .type(TbMsgType.POST_TELEMETRY_REQUEST) + .originator(customer.getId()) + .metaData(TbMsgMetaData.EMPTY) + .data(TbMsg.EMPTY_JSON_OBJECT) + .build(); + + var expectedMsg = msg.transform() + .originator(CUSTOMER_ID) + .build(); + given(ctxMock.transformMsgOriginator(any(TbMsg.class), any(EntityId.class))).willReturn(expectedMsg); node.init(ctxMock, new TbNodeConfiguration(JacksonUtil.valueToTree(config))); + + // WHEN node.onMsg(ctxMock, msg); - then(deviceServiceMock).should().findDeviceById(TENANT_ID, DEVICE_ID); + // THEN then(ctxMock).should().transformMsgOriginator(msg, CUSTOMER_ID); ArgumentCaptor actualMsg = ArgumentCaptor.forClass(TbMsg.class); then(ctxMock).should().tellSuccess(actualMsg.capture()); @@ -216,8 +270,6 @@ public class TbChangeOriginatorNodeTest { .originator(TENANT_ID) .build(); - given(ctxMock.getDbCallbackExecutor()).willReturn(dbExecutor); - given(ctxMock.getTenantId()).willReturn(TENANT_ID); given(ctxMock.transformMsgOriginator(any(TbMsg.class), any(EntityId.class))).willReturn(expectedMsg); node.init(ctxMock, new TbNodeConfiguration(JacksonUtil.valueToTree(config))); @@ -240,9 +292,6 @@ public class TbChangeOriginatorNodeTest { .data(TbMsg.EMPTY_JSON_OBJECT) .build(); - given(ctxMock.getDbCallbackExecutor()).willReturn(dbExecutor); - given(ctxMock.getRelationService()).willReturn(relationServiceMock); - given(ctxMock.getTenantId()).willReturn(TENANT_ID); given(relationServiceMock.findByQuery(any(TenantId.class), any(EntityRelationsQuery.class))).willReturn(Futures.immediateFuture(Collections.emptyList())); node.init(ctxMock, new TbNodeConfiguration(JacksonUtil.valueToTree(config))); @@ -282,9 +331,6 @@ public class TbChangeOriginatorNodeTest { .originator(DEVICE_ID) .build(); - given(ctxMock.getDbCallbackExecutor()).willReturn(dbExecutor); - given(ctxMock.getAlarmService()).willReturn(alarmServiceMock); - given(ctxMock.getTenantId()).willReturn(TENANT_ID); given(alarmServiceMock.findAlarmByIdAsync(any(TenantId.class), any(AlarmId.class))).willReturn(Futures.immediateFuture(alarm)); given(ctxMock.transformMsgOriginator(any(TbMsg.class), any(EntityId.class))).willReturn(expectedMsg); @@ -315,9 +361,6 @@ public class TbChangeOriginatorNodeTest { .originator(ASSET_ID) .build(); - given(ctxMock.getDbCallbackExecutor()).willReturn(dbExecutor); - given(ctxMock.getAssetService()).willReturn(assetServiceMock); - given(ctxMock.getTenantId()).willReturn(TENANT_ID); given(assetServiceMock.findAssetByTenantIdAndName(any(TenantId.class), any(String.class))).willReturn(new Asset(ASSET_ID)); given(ctxMock.transformMsgOriginator(any(TbMsg.class), any(EntityId.class))).willReturn(expectedMsg); @@ -355,9 +398,6 @@ public class TbChangeOriginatorNodeTest { .data(TbMsg.EMPTY_JSON_OBJECT) .build(); - given(ctxMock.getDbCallbackExecutor()).willReturn(dbExecutor); - given(ctxMock.getAssetService()).willReturn(assetServiceMock); - given(ctxMock.getTenantId()).willReturn(TENANT_ID); given(assetServiceMock.findAssetByTenantIdAndName(any(TenantId.class), any(String.class))).willReturn(null); node.init(ctxMock, new TbNodeConfiguration(JacksonUtil.valueToTree(config))); 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 deleted file mode 100644 index 0be4e64ca2..0000000000 --- a/rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/util/EntitiesCustomerIdAsyncLoaderTest.java +++ /dev/null @@ -1,154 +0,0 @@ -/** - * Copyright © 2016-2025 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 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.TestDbCallbackExecutor; -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.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 SUPPORTED_ENTITY_TYPES = EnumSet.of( - EntityType.CUSTOMER, - EntityType.USER, - EntityType.ASSET, - EntityType.DEVICE - ); - private static final ListeningExecutor DB_EXECUTOR = new TestDbCallbackExecutor(); - @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(device).when(deviceServiceMock).findDeviceById(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()); - } - } - } - -}